项目名BS-IM时间2023.06-2023.10 上海百胜软件股份有限公司 Java开发工程师组长项目描述BS-IM是一款基于高性能Netty框架设计的百万级聊天消息系统采用Spring cloud Gateway、Redis 和Netty中间件基于事件驱动的方式进行消息通信结合Kafka消息队列向多数据库存储如Mysql和ClickHouse等。基于中台思想抽取出了新零售行业的消息通信的共性能力向上拓展即可实现智能客服、多商家入驻订阅与消息推送和直播带货万人群聊等功能同时支持弹性伸缩的方式满足企业的实际业务需求。业务难点1.心跳机制与断线重连难点。2.海量聊天数据存储难点。3.高性能离线消息读取难点。4.百万级消息订阅与推送难点。5.百万级在线直播难点。6.消息发送ACK机制难点。7.消息未读数统计难题。实现流程1.启动Netty同时向Zookeeper注册服务地址采用临时顺序节点方式GateWay进行Watch节点变化。2.智能客服U2发起连接调用至GateWay通过自定义负载均衡算法选择一台Netty建立起Channel连接。并且记录至Redis计。客户端U1同理与NettyServer建立连接。3.U1查询关系服务查询到智能客服的UID为U2发送消息U1-U2至GateWay从Redis取得U2的Netty地址N2将消息转给N2N2通过U2的Channel发送消息至U2客户端。同时GateWay将消息发送至Kafka。4.聊天服务消费Kafka消息至存储Redis基于淘汰的最新热点消息消息的索引号、ClickHouseHash分区存储全量消息、ElasticSearch建立索引用于消息检索。5.在离线场景时客户端上线优先通过离线服务分页拉取Redis的最新离线消息以及消息状态数当客户端查看到消息后发送ACK至存储服务更新Redis已读标志更多历史消息从ClickHouse获取定期做数据归档策略。6.利用Netty读空闲机制建立客户端与服务器端心跳检测。若Server宕机GateWay利用Watch机制除去服务列表地址等待客户端感知自动多次重连更新原有Redis服务地址键值对。若Client宕机则通过Server去除Redis服务地址键值对。介绍一下百胜im就是客服聊天系统。这个是专门为电商领域开发的中台聊天系统。业务背景是当时是丝芙兰客户还没签约其中一点就是他们需要一个智能客服。参考我们团队当时参考了飞猪他们内部的im架构简化了架构模型结合我们的电商场景做最通用的下沉能力的中台方便后续客户进行敏捷二开。做了自己的架构方案智能预料库由于我们没ai团队用了别家公司的数睿调用他们的ai智能回答的一个Pythonweb的包。一个前端、后端就是我全身心投入除了违规词检测和es检索消息后端都是我写的安卓、ios和web都是uniapp一套一个前端开发阶段大概花了三个月一个半月的测试修bug。我按照服务角度讲起首先是服务层的关系服务这个是聊天关系的查询服务用户通过查询自己的关联的聊天对象的userid不管是人工客服、智能客服、商家都有自己的userid接着就是已自己的userid发送消息给对方的userid每个userid其实就对应一个client首先消息方请求网关验证token之后会根据负载均衡算法去连接一台netty server建立起长连接并且将userid 和 几号 netty server存储到 redis里。同样消息接受方也已这种方式去和自己的netty server建立长连接存储到redis里。接着发送消息给消息接受方消息到网关通过查询redis就知道接受方和那一台netty server保持着channle通道就把消息转发给具体的netty server。接着netty server就把消息转给消息接受方了发送成功。netty server 启动的时候会将自己的ip 和端口信息以及临时顺序节点的方式挂载到zk的一个父节点下面。之后网关会去读取这个父节点下面的全部节点在自己内部维护全部的netty server地址并且watch机制这个父节点的子节点是否变化。若netty server 下线因为是临时顺序节点所有zk感知到以后节点会移除并且通知到了GatewayGateway就会在自己的服务列表里面移除对应的netty server地址。由于clinet 和 server 是长连接下线是可以感知到的clinet就会去网关里面删除之前保存在redis的连接信息并且通过网关维护的最新地址建立新的连接。若clinet端下线了server感知到以后会去redis删除该userid的连接信息。10秒C会S发生ping pong 包采用读空闲机制15秒未读到数据则判断上一次的心跳包和当前时间是否大于30S大于则判断对方下线。由于client在网外环境发生网络瞬断或者隔离了这种情况长连接是感知不到的所以需要让client 和sever 即在不进行任何数据传送的情况下定时发送心跳包目标是发生这种网络断了也要移除对应的客户端。客户端会采用netty写空闲的方式就是当客户端5秒内都没数据将数据写入通道就会发送一个ping的文字信条数据服务端接受以后会回一个pong的文字信息并且将该心跳包的创建时间写入chanle。服务端会采用读空闲的机制若20秒没有读取到客户端的数据则查询通道内上一次心跳包的创建时间比较当前时间是否大于30秒若是则进行连接关闭出现了网络隔离。客户端 服务端若客户端一直写数据到服务端空闲不会触发网络不隔离。若服务端一直写数据到客户端空闲触发检测心跳包。负载均衡算法拿了ribbon的代码轮询加1取模里了拓展接口。令牌桶限流为了以后的有直播场景的万人群聊场景创建一个令牌桶主播只需要一个令牌网关就能过进行消息群播没消费的那些用户可能需要20个令牌这样进行令牌桶限流人多的时候要进行按照用户等级进行筛选发言人少的时候大家都能发言也就是说某些人的发言只有自己设备上看得到并没有进行群播并且能抵御突发流量。违规词根据睿数提供的直播和电商场景的违规词构建了一颗trie树前缀树同时将进来的消息进行分词进行前缀算法匹配。接着是最复杂的存储问题。我们第一个实践客户用的就是单机mysql。所有消息都会发kafka让mysql存储同时用户已读以后会回调存储进行消息状态修改为已读。用户离线以后会进行离线未读消息的分页批量拉取。就这样交付了之后呢上线一个多月拉取离线消息就变得很慢客户觉得卡了因为到了mysql的四层b树之后我们就进行了不停机、服务无缝切换数据源至分库分表的数据源的实践采用chanl 配合 mq进行。根据每条消息的发送方id进行分区可满足我们全部的查询条件因为分库分最大的难点在于分片键的选择不然根据其他不分片的键查询就得全库表扫描使用sprding-jdbc框架进行客户端分4库8表 取模32。接下了数据。同时上线了一个离线服务将活跃用户的最新的离线消息缓存至了redis避免活跃用户频繁上下线查询数据库并且作为补偿增加了未读消息数的功能。用户的离线收件箱sendid {3 [1,2,3]}未读数量消息id号userid1 哈哈2 你好消息主表消息id 发送人id 接收人id 消息内容 消息状态1 1 2 msg 未读消息索引表参与者1id 参与者2id 消息id 接受类型1 2 1 发送2 1 1 接受1 和 2 消息记录1手机里面对2的聊天记录owner_id other_id mid box_type1 2 1 1(收)1 2 2 2(发)1 2 3 1(收)Redishash结构大key hashkey hashvalue1 2 11 3 6读写扩散即群聊。userid userid userid1 2 3gid1userid 1 的群聊收件箱参与1 参与2 收件 聊消息号1 1 1 11 1 1 2写扩散分两种一种群聊就是普通商家粉丝就几百那么就当做几百个人和这个商家建立一个群只不过是只有商家可以发言和单聊一样发送消息全部发送发到每个粉丝的收件箱写扩散机制。还有一种是品牌方订阅就是全部用户都能看到该消息那么就是品牌方的消息只在自己的发件箱里等到用户上线的时候去查询自己的关系服务里面的品牌订阅方的拉取它的发送过的最新消息到自己的收件箱读扩散机制。不停机数据迁移业务背景是我们的消息系统上线了原始版本是单机mysql后来发现用户的离线消息批量查询不动了之后我们决定进行在线同步数据不停机至分库分表的数据库。不停机数据源数据进行同步最主要担心的点在于当主库的数据已经同步到目标库了此时主库的某一行记录被修改了也必须要在目标库同步数据。所以我们采用定时任务批量同步历史数据对于同步开始时动态变化的数据使用canal 监听binlong的方式通过唯一区分键进行幂等。1.批量查询源数据表数据发往kafka下游消费者监听进行hash分区入库于此同时开启binglog监听源数据库发往kafka不进行下游消费。2.等待离线数据同步任务的数据结束开启动态binlog的消费者。insert 时 有唯一的id update 可以拿到原始的id 数据更新后的数据对于新记录id查重入库对于update 根据id无脑更新这样子的话当历史数据已经同步过去了也会随着binglog更新进行更新。
