简介:这份文档面向后端架构师、大数据开发工程师及对实时计算感兴趣的技术人员,系统梳理了携程实时用户行为服务从旧架构痛点出发的完整重构实践。内容围绕推举系统、动态广告、用户画像、浏览历史等真实业务场景,剖析数据覆盖不全、输出格式不统一、日志模块性能不足等问题,并给出处理流与输出流的双流设计方案。技术选型部分详解Java、Kafka、Storm、Redis、MySQL、Tomcat与Spring的取舍依据,随后从实时性、可用性、功能性、扩展性四个维度展开,涵盖突发流量洪峰应对、双队列补偿重试、积压数据消解、全栈集群化与DB降级等关键设计,并附有系统日处理约20亿数据、上线可用约300毫秒、查询平均延迟约6毫秒的实测指标。资源为1个docx文档,压缩包约161KB,已有130人学习,适合用于架构方案参考与面试复盘。
1. 从 20 亿条行为日志说起:这套架构到底解决了什么
每天 20 亿条用户行为日志,从 App、H5、Online 三端上报,到能被推荐系统、动态广告、用户画像、浏览历史这些下游场景查到,中间只有 300 毫秒。查询服务每天扛 8000 万次请求,平均延迟 6 毫秒。这不是某个实验室的 benchmark,是携程实时用户行为服务系统在生产环境跑出来的数字。
这套东西的本质,是一个把「用户刚点了什么」变成「下游立刻能用」的基础服务。猜你喜欢要拿它做实时推荐,广告系统要拿它做动态投放,用户画像要拿它做实时更新。它不直接面向 C 端用户,但 C 端每一次刷出新推荐、每一次看到广告变化,背后都有它在跑。
适合谁看?如果你正在设计或维护一套实时数据链路,数据量在千万到十亿级别,下游有多个业务方要接,同时对延迟和可用性有硬要求,那这份架构实践里的取舍和踩坑记录,基本可以当参考模板用。如果你只是做个小规模埋点统计,这里面的双队列、降级、分片扩容可能偏重,但思路仍然值得看一遍。
2. 技术栈选型:为什么是 Java+Kafka+Storm+Redis+MySQL
2.1 选型背后的真实约束
这套系统 2021 年落地,技术栈是 Java + Kafka + Storm + Redis + MySQL + Tomcat + Spring。看起来都是「老面孔」,但每一个选择背后都有具体约束,不是拍脑袋定的。
Java 是公司内部氛围决定的,更关键的是 Java 生态里大数据组件成熟,Storm、Kafka 的客户端都是 Java 原生,团队维护成本低。Kafka 作为分布式消息队列在公司内部已经有成熟应用,运维支持环境现成,不需要重新造轮子。Storm 作为流计算框架也已经落地,有运维支持,能快速上线。
Redis 被选中的原因是它的 HA、SortedSet 和过期特性。SortedSet 可以按时间戳排序,天然适合存用户行为序列;过期特性可以自动清理冷数据,不用额外写清理任务。MySQL 的选择更有意思——对比 HBase 和 ElasticSearch,在十亿数据级别上,MySQL 的稳定性和功能表现更好,而且经过水平切分设计后,水平扩展能力并不差。
提示:选型时不要只看「哪个技术新」,要看「哪个技术在这个数据量级、这个团队背景下,运维成本最低、出问题最好查」。
2.2 数据流向与模块职责
系统有两条数据流:处理流和输出流。
处理流:客户端(App/Online/H5)上传行为日志到 CollectorService,CollectorService 把消息发到 Kafka,Storm 从 Kafka 读数据,处理之后写入数据层(Redis + MySQL)。
输出流:Web Service 后台从数据层拉数据,输出给调用方。内部服务调用比如推荐系统,前台输出比如浏览历史。
用代码块表示一下核心消费逻辑的结构:
// Storm Bolt 中消费 Kafka 消息并写入 Redis + MySQL 的核心逻辑 // 注意:这是结构示意,不是完整可运行代码 public class BehaviorProcessBolt extends BaseRichBolt { private JedisCluster redisCluster; private DataSource mysqlDataSource; private KafkaProducer<String, String> retryProducer; @Override public void execute(Tuple tuple) { String message = tuple.getStringByField("value"); try { UserBehavior behavior = parse(message); // 写入 Redis,按用户 ID 分片,用 SortedSet 按时间排序 redisCluster.zadd("behavior:" + behavior.getUserId(), behavior.getTimestamp(), behavior.toJson()); // 写入 MySQL,按用户 ID 水平切分 writeToMySQL(behavior); } catch (Exception e) { // 写入失败,转入重试队列,不阻塞主队列消费 retryProducer.send(new ProducerRecord<>("behavior-retry", message)); } } }逻辑说明:这段代码的核心是「主流程写 Redis + MySQL,失败转重试队列」。参数上,Redis 的 key 按用户 ID 分片,value 用 SortedSet 按时间戳排序,方便按时间范围拉取。MySQL 写入按用户 ID 水平切分,分片数量选 2 的 n 次方,为后续扩容留余地。重试队列是独立的 Kafka topic,不阻塞主队列消费。
2.3 双队列设计:怎么保证新数据不被旧问题拖死
实时系统最怕的不是数据多,是「一条坏数据把整条链路堵死」。比如数据库连接超时,如果处理程序一直等,后面新来的数据就全积压了。
这套系统用了双队列设计。生产者把行为记录写入 Queue1,Worker 从 Queue1 消费新数据。如果遇到异常数据(比如数据库连不上),Worker 把异常数据写入 Queue2,自己继续消费 Queue1 的新数据。RetryWorker 监听 Queue2,按策略等待或重新写入 Queue2,直到处理成功。
# 双队列的 Kafka topic 配置示意 # Queue1:主队列,保持数据新鲜度 kafka-topics.sh --create --topic behavior-main \ --partitions 32 --replication-factor 3 \ --config retention.ms=86400000 # Queue2:重试队列,存放异常数据 kafka-topics.sh --create --topic behavior-retry \ --partitions 16 --replication-factor 3 \ --config retention.ms=604800000参数说明:主队列 retention 设 1 天,保证新数据优先;重试队列 retention 设 7 天,给异常数据足够的重试窗口。分区数主队列 32、重试队列 16,按流量比例分配。副本数都是 3,保证可用性。
注意:双队列的关键是「Worker 对 Queue1 的消费进度不被 Queue2 影响」。如果 RetryWorker 处理太慢导致 Queue2 积压,不会反过来阻塞主队列。
3. 实时性与可用性:Storm 的 at least once 和降级开关怎么配合
3.1 为什么选 at least once 而不是 exactly once
Storm 支持三种消息保证策略:at least once、at most once、exactly once。实时用户行为系统选的是 at least once。
原因很直接:对用户行为数据来说,首要目标是「尽量少丢」,而不是「绝对不重」。exactly once 需要事务支持,会降低吞吐量,而且实现复杂度高。at least once 允许消息重发,所以程序处理必须实现幂等——同一条行为记录重复写入,结果要一样。
幂等实现常见做法是:用行为日志里的唯一 ID(比如 requestId + 时间戳)做去重键,写入 Redis 时用 SETNX 或者 SortedSet 的 score 去重,写入 MySQL 时用 INSERT IGNORE 或者 ON DUPLICATE KEY UPDATE。
-- MySQL 幂等写入:用唯一索引 + INSERT IGNORE 实现 -- 表结构:behavior_log,唯一键是 (user_id, behavior_id) CREATE TABLE behavior_log ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id BIGINT NOT NULL, behavior_id VARCHAR(64) NOT NULL, behavior_type VARCHAR(32), timestamp BIGINT, extra JSON, UNIQUE KEY uk_user_behavior (user_id, behavior_id) ) ENGINE=InnoDB; -- 写入时用 INSERT IGNORE,重复数据自动跳过 INSERT IGNORE INTO behavior_log (user_id, behavior_id, behavior_type, timestamp, extra) VALUES (?, ?, ?, ?, ?);逻辑说明:唯一索引uk_user_behavior保证同一用户的同一条行为只存一次。INSERT IGNORE在遇到重复键时直接跳过,不报错。这样即使 Storm 重发消息,MySQL 里也不会出现重复数据。
3.2 突发流量洪峰怎么扛
Storm 的 scale out 能力是应对流量洪峰的核心。通过后台修改 worker 数量参数,重启 topology,就能扩展计算能力。不需要改代码,不需要重新打包,只是调整并行度。
# Storm topology 并行度调整示意 # 提交 topology 时指定 worker 数量 storm jar behavior-service.jar com.ctrip.behavior.BehaviorTopology \ behavior-topology \ -c topology.workers=8 \ -c topology.acker.executors=8 # 运行时动态调整(部分 Storm 版本支持) storm rebalance behavior-topology -w 16 -n 4参数说明:topology.workers是 worker 进程数,每个 worker 是一个 JVM。topology.acker.executors是 acker 线程数,负责消息确认。rebalance命令可以在不重启 topology 的情况下调整并行度,但部分版本支持有限,常见做法还是改配置后重启。
提示:重启 topology 时,Kafka 记录的消费游标会保留,程序重启后从上次位置继续消费,不会丢数据。但要注意,如果重启期间有大量数据积压,需要评估消费速度能否追上。
3.3 降级开关:DB 挂了怎么办
系统可用性设计里,降级是最后一道防线。正常流程是 Storm 从 Kafka 读数据,分别写入 Redis 和 MySQL。服务从 Redis 拉数据,取不到时从 DB 补偿。
当 MySQL 不可用时,打开 DB 降级开关:Storm 正常写 Redis,但不再写 MySQL。数据进 Redis 就能被查询服务使用。同时 Storm 把数据写一份到 Kafka 的 retry 队列。MySQL 恢复后,关闭降级开关,Storm 消费 retry 队列,把数据补写入 MySQL。
// 降级开关的简单实现:用配置中心或 Redis 标志位控制 public class DegradeSwitch { private static volatile boolean dbDegrade = false; private static volatile boolean redisDegrade = false; // 配置中心推送或定时轮询更新开关状态 public static void updateSwitch(String key, boolean value) { if ("db.degrade".equals(key)) { dbDegrade = value; } else if ("redis.degrade".equals(key)) { redisDegrade = value; } } public static boolean isDbDegrade() { return dbDegrade; } }逻辑说明:降级开关用 volatile 保证多线程可见性,通过配置中心推送或定时轮询更新。Storm Bolt 在写入前检查开关状态,决定是否跳过 MySQL 写入。Redis 降级类似,但 Redis 服务能力远超过 MySQL,降级时吞吐量下降,需要监控 DB 压力,必要时临时停止数据写入。
注意:降级期间 Redis 和 MySQL 数据会不一致,但系统恢复后通过 retry 队列保证最终一致性。这个「最终一致」的时间窗口取决于 retry 队列的消费速度,需要提前评估。
4. 扩展性与部署:MySQL 分片扩容和 Storm 多版本运行
4.1 MySQL 水平切分:分片数为什么选 2 的 n 次方
系统要求支撑 10 倍容量扩展,最难的部分在数据层,因为涉及存量数据迁移。Redis 实现了一致性哈希,扩容时加机器、对新分区数据做读补偿就行。MySQL 做了水平切分,分片数量选 2 的 n 次方。
为什么是 2 的 n 次方?因为携程 MySQL 普遍是一主一备部署。扩容时可以直接把备机拉平成第二台主机。假设原来分了 2 个库 d0 和 d1,都放在服务器 s0 上,s0 有备机 s1。扩容步骤:
- 确保 s0 -> s1 同步顺利,没有明显延迟
- s0 临时关闭读写权限
- 确认 s1 已经完全同步 s0 更新
- s1 开放读写权限
- d1 的 DNS 由 s0 切换到 s1
- s0 开放读写权限
整个过程利用 MySQL 复制分发特性,避免人工同步,几分钟完成。结合 DB 降级功能,只在 DNS 切换的几秒钟产生异常。
-- 分片路由示意:按 user_id 取模路由到不同库 -- 分片数 shardCount = 4(2 的 2 次方) -- 扩容时 shardCount 变为 8,只需迁移一半数据 public class ShardRouter { private static final int SHARD_COUNT = 4; public static String getDataSourceKey(long userId) { int shard = (int) (userId % SHARD_COUNT); return "ds_" + shard; } }参数说明:SHARD_COUNT是分片数,选 2 的 n 次方。扩容时从 4 变 8,只需要把原来每个分片的数据拆一半到新分片,迁移量可控。如果选 3 或者 5,扩容时数据迁移会复杂很多。
4.2 Storm 部署与多版本运行
Storm 部署很简单:上传更新程序 jar 包,重启任务。部署后程序上下文丢失,但可以通过 Kafka 记录的游标找到之前处理位置,恢复处理。
有些情况下程序需要多版本运行,比如行为记录临时有多个版本。这时新增一个 backupJob,在 backupJob 中运行历史版本。主 Job 处理新版本数据,backupJob 处理旧版本数据,两者互不干扰。
# 提交主 Job 和 backupJob 的示意 # 主 Job:处理当前版本行为数据 storm jar behavior-service.jar com.ctrip.behavior.MainTopology \ behavior-main-topology # backupJob:处理历史版本行为数据 storm jar behavior-service-v1.jar com.ctrip.behavior.BackupTopology \ behavior-backup-topology逻辑说明:两个 topology 消费不同的 Kafka topic 或同一 topic 的不同分区。主 Job 用新版本代码,backupJob 用旧版本代码。这样在数据格式过渡期间,新旧数据都能被正确处理。
提示:多版本运行会增加运维复杂度,常见做法是尽量在代码里做版本兼容,而不是长期维护两个 Job。backupJob 只作为过渡方案,数据格式统一后及时下线。
4.3 积压数据消解:调整消费游标和 backupWorker
数据积压时,可以调整 Worker 的消费游标,从最新数据重新开始消费,保证最新数据得到处理。两头未处理的一段数据,启动 backupWorker,指定起止游标,消费完指定区间后自动停止。
# 调整 Kafka 消费游标到最新位置(示意) # 使用 kafka-consumer-groups 工具 kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group behavior-worker-group \ --topic behavior-main \ --reset-offsets --to-latest --execute # backupWorker 指定起止游标消费(伪代码示意) # backupWorker.setStartOffset(1000000); # backupWorker.setEndOffset(2000000); # backupWorker.run(); // 消费完自动停止参数说明:--reset-offsets --to-latest把消费组偏移量重置到最新,Worker 从最新数据开始消费。backupWorker 需要自己实现起止游标控制,消费到 endOffset 后自动停止。这样既保证了新数据不延迟,又不会丢掉积压的数据。
5. 避坑与排查:这套架构落地时最容易翻车的五个点
5.1 幂等没做全,重试导致数据重复
现象:Storm 重发消息后,MySQL 里出现重复行为记录,下游统计指标偏高。
原因:at least once 策略下消息可能重发,如果写入逻辑没有幂等保证,重复消息就会产生重复数据。
解决:用行为日志的唯一 ID 做去重键,MySQL 建唯一索引,写入用 INSERT IGNORE 或 ON DUPLICATE KEY UPDATE。Redis 写入用 SETNX 或 SortedSet score 去重。幂等要覆盖所有写入路径,不能只做 MySQL 不做 Redis。
5.2 降级开关忘记关,数据长时间不一致
现象:MySQL 恢复后,Redis 和 MySQL 数据长时间不一致,下游查询结果忽新忽旧。
原因:DB 降级开关打开后忘记关闭,Storm 一直不写 MySQL,retry 队列持续积压。
解决:降级开关加自动超时机制,比如打开后 30 分钟自动尝试恢复。同时加监控告警,降级开关打开超过阈值就通知值班人员。恢复后要验证 retry 队列消费进度,确认数据补写完成。
5.3 MySQL 分片扩容时 DNS 切换产生异常
现象:扩容切换 DNS 的几秒钟内,部分查询请求失败或超时。
原因:DNS 切换有传播延迟,客户端缓存了旧 DNS 记录,切换瞬间连接不上。
解决:结合 DB 降级功能,在 DNS 切换前打开降级开关,让查询走 Redis。切换完成后关闭降级开关。另外 DNS TTL 设短一点,比如 60 秒,减少传播延迟。
5.4 Storm 重启后消费游标丢失,数据重复消费
现象:Storm 重启后,从 Kafka 最早位置开始消费,大量数据重复处理。
原因:Kafka 消费游标没有正确提交,或者 Storm 的 spout 没有配置从上次位置恢复。
解决:确保 Kafka spout 配置了正确的 consumer group 和 auto.offset.reset 策略。常见做法是设 auto.offset.reset=latest,重启后从最新位置消费。如果需要精确恢复,用 Kafka 的 offset 管理 API 手动提交和恢复。
5.5 Redis 降级时吞吐量下降,MySQL 压力过大
现象:Redis 降级后,查询请求全部打到 MySQL,MySQL 压力飙升,响应变慢。
原因:Redis 服务能力远超过 MySQL,降级后流量全部转移到 MySQL,超出其承载能力。
解决:Redis 降级时监控 MySQL 压力,如果压力过大,临时停止数据写入,降低 MySQL 负载,优先保证查询服务稳定。同时评估是否需要对 MySQL 做限流或排队。
6. 从 300 毫秒到 6 毫秒:验证这套架构是否真的跑通了
这套架构最终跑出来的数字是:每天处理 20 亿条数据,数据从上线到可用 300 毫秒左右,查询服务每天 8000 万次请求,平均延迟 6 毫秒。怎么验证你的系统也能达到类似水平?我一般会从三个维度做验证。
第一,端到端延迟验证。在客户端埋一个时间戳,在查询服务输出时再打一个时间戳,两个时间戳的差值就是端到端延迟。采样 1% 的请求,统计 P50、P95、P99。如果 P99 超过 500 毫秒,说明链路中有瓶颈,需要逐段排查。
第二,数据一致性验证。在降级恢复后,随机抽样对比 Redis 和 MySQL 的数据,确认 retry 队列消费完成后两边数据一致。常见做法是写一个对账脚本,按用户 ID 抽样,对比 Redis 和 MySQL 的行为记录条数和内容。
# 对账脚本示意:抽样对比 Redis 和 MySQL 数据 import redis import pymysql import random r = redis.Redis(host='redis-host', port=6379) conn = pymysql.connect(host='mysql-host', user='user', password='pass', db='behavior') def check_consistency(user_id): # 从 Redis 拉取行为记录 redis_data = r.zrange(f"behavior:{user_id}", 0, -1) # 从 MySQL 拉取行为记录 with conn.cursor() as cur: cur.execute("SELECT behavior_id FROM behavior_log WHERE user_id=%s", (user_id,)) mysql_data = [row[0] for row in cur.fetchall()] # 对比条数和内容 if len(redis_data) != len(mysql_data): print(f"User {user_id}: Redis={len(redis_data)}, MySQL={len(mysql_data)}") return False return True # 随机抽样 1000 个用户 for _ in range(1000): uid = random.randint(1, 10000000) check_consistency(uid)逻辑说明:这个脚本随机抽样用户,对比 Redis 和 MySQL 的行为记录条数。如果条数不一致,说明降级恢复后数据没有完全同步。参数上,抽样比例根据数据量调整,数据量大时抽 0.1% 就够。
第三,压力测试验证。用压测工具模拟突发流量,观察 Storm 的 scale out 是否及时,Kafka 是否积压,Redis 和 MySQL 的响应时间是否在可接受范围。常见做法是用 JMeter 或 wrk 打查询服务,同时用 Kafka 生产者灌数据,观察端到端延迟变化。
注意:压测时不要只打查询服务,要同时灌数据。只打查询不灌数据,测不出处理流的瓶颈。只灌数据不打查询,测不出输出流的瓶颈。
从那以后我每次设计实时链路,都会强制走一遍「端到端延迟 + 数据一致性 + 压力测试」这三步。少一步,上线后都可能出玄学问题。希望帮到你。
本文还有配套的精品资源,点击获取