
1. 项目概述在社交平台和内容社区中用户关系链关注/粉丝系统是最核心的基础设施之一。一个典型的中大型平台每天会产生数百万甚至上千万的关注关系变更操作这对系统的实时性、一致性和可用性提出了极高要求。传统基于数据库直接读写的方式在流量高峰时经常出现性能瓶颈而简单的缓存方案又难以保证数据强一致性。我们设计的这套基于Canal Kafka的高可用关注系统采用了一主多从的架构模式通过MySQL Binlog日志解析、消息队列异步处理和分布式缓存的多层配合实现了每秒10万级关系变更操作的稳定处理。系统核心指标包括数据延迟控制在200ms内主从切换时数据丢失窗口1秒99.9%的请求响应时间在50ms以下2. 核心架构设计2.1 技术选型解析Canal组件选择 作为阿里巴巴开源的MySQL binlog增量订阅组件Canal相比其他方案如Debezium具有以下优势对中文文档支持完善社区活跃度高原生支持Kafka协议输出经过双十一等大促场景验证提供完善的管理控制台我们在生产环境使用Canal 1.1.6版本主要配置参数包括canal.instance.mysql.slaveId 1234 canal.mq.topicuser_relation canal.mq.partition0Kafka集群配置 采用3 broker组成的集群关键配置如下消息保留策略log.retention.hours72分区数num.partitions6与关系类型数量匹配副本因子default.replication.factor2特别注意必须开启Kafka的acksall配置确保消息不丢失。我们曾因未设置此参数导致主从切换时部分关注关系丢失。2.2 一主多从关系链实现系统架构如下图所示文字描述主MySQL集群处理所有写操作配置半同步复制从MySQL集群3个只读实例通过Proxysql实现负载均衡Canal Server部署在独立节点伪装为MySQL从库Kafka集群接收Canal的binlog事件处理服务集群消费Kafka消息更新Redis和ES关系链的存储采用正反向双写策略正向关系关注列表user_id - [followee_id1, followee_id2...]反向关系粉丝列表user_id - [follower_id1, follower_id2...]3. 关键实现细节3.1 数据同步流程完整的数据流转过程如下用户在客户端执行关注操作API服务写入主MySQL的relation表Canal捕获binlog事件并发送到Kafka处理服务消费消息并执行public void process(Message message) { // 更新Redis缓存 redisTemplate.opsForSet().add( follow: message.getUserId(), message.getFolloweeId()); // 更新ES索引 esClient.updateRelation( message.getUserId(), message.getFolloweeId(), message.getOperationType()); }从库通过Proxysql提供读服务3.2 高可用保障措施主从切换方案 使用Orchestrator工具实现自动故障转移关键配置DetectClusterAliasQuery: SELECT user_relation PromotionIgnoreHostnameFilters: standby.mysql数据一致性校验 每天凌晨通过校验任务比对MySQL和Redis的数据差异-- 校验SQL示例 SELECT user_id, COUNT(*) as db_count FROM user_relations GROUP BY user_id消息幂等处理 在Kafka消费者端实现去重逻辑def handle_message(msg): if redis.get(fmsg:{msg.id}): return # 处理逻辑... redis.setex(fmsg:{msg.id}, 3600, 1)4. 性能优化实践4.1 缓存设计技巧采用分级缓存策略一级缓存本地Caffeine缓存最近活跃用户Caffeine.newBuilder() .maximumSize(10_000) .expireAfterWrite(5, TimeUnit.MINUTES) .build();二级缓存Redis集群使用Hash结构存储关系数据设置合理的过期时间通常24小时4.2 Kafka调优经验分区策略优化按用户ID哈希分配到不同分区确保同一用户的关系变更顺序处理消费者配置fetch.min.bytes65536 fetch.max.wait.ms100 max.poll.records500监控指标消费延迟kafka.consumer.lag处理耗时process.time5. 典型问题排查5.1 数据同步延迟现象用户反映关注状态不同步排查步骤检查Canal位点是否落后curl http://canal:11111/destinations/user_relation/metrics查看Kafka堆积情况kafka-consumer-groups --describe \ --bootstrap-server kafka:9092 \ --group relation_consumer检查处理服务监控指标解决方案增加处理服务实例数调整消费者max.poll.records参数优化Redis批量操作5.2 主从切换异常现象切换后部分关系丢失原因Kafka未配置acksall导致消息丢失修复方案重新配置Kafka生产者acksall retries3实现补偿查询机制public void checkConsistency(long userId) { // 对比DB和缓存差异 }6. 扩展与演进当前系统支持的功能扩展方向关系变更通知通过Kafka Connect将数据同步到通知服务实现提及的实时推送图数据库集成将关系数据同步到Neo4j支持多层关系分析冷热数据分离将3个月前的关注关系归档到HBase减少主库数据量这套架构经过618大促验证峰值时成功处理了15万QPS的关注操作。在实际部署时建议至少准备3台16核32G的物理机运行Kafka集群MySQL配置建议innodb_buffer_pool_size为总内存的70%。