ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

20亿日志300毫秒可见:携程实时用户行为系统架构实践

20亿日志300毫秒可见:携程实时用户行为系统架构实践 简介这份文档面向后端架构师、大数据开发工程师及对实时计算感兴趣的技术人员系统梳理了携程实时用户行为服务从旧架构痛点出发的完整重构实践。内容围绕推举系统、动态广告、用户画像、浏览历史等真实业务场景剖析数据覆盖不全、输出格式不统一、日志模块性能不足等问题并给出处理流与输出流的双流设计方案。技术选型部分详解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. 技术栈选型为什么是 JavaKafkaStormRedisMySQL2.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上传行为日志到 CollectorServiceCollectorService 把消息发到 KafkaStorm 从 Kafka 读数据处理之后写入数据层Redis MySQL。输出流Web Service 后台从数据层拉数据输出给调用方。内部服务调用比如推荐系统前台输出比如浏览历史。用代码块表示一下核心消费逻辑的结构// Storm Bolt 中消费 Kafka 消息并写入 Redis MySQL 的核心逻辑 // 注意这是结构示意不是完整可运行代码 public class BehaviorProcessBolt extends BaseRichBolt { private JedisCluster redisCluster; private DataSource mysqlDataSource; private KafkaProducerString, 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 双队列设计怎么保证新数据不被旧问题拖死实时系统最怕的不是数据多是「一条坏数据把整条链路堵死」。比如数据库连接超时如果处理程序一直等后面新来的数据就全积压了。这套系统用了双队列设计。生产者把行为记录写入 Queue1Worker 从 Queue1 消费新数据。如果遇到异常数据比如数据库连不上Worker 把异常数据写入 Queue2自己继续消费 Queue1 的新数据。RetryWorker 监听 Queue2按策略等待或重新写入 Queue2直到处理成功。# 双队列的 Kafka topic 配置示意 # Queue1主队列保持数据新鲜度 kafka-topics.sh --create --topic behavior-main \ --partitions 32 --replication-factor 3 \ --config retention.ms86400000 # Queue2重试队列存放异常数据 kafka-topics.sh --create --topic behavior-retry \ --partitions 16 --replication-factor 3 \ --config retention.ms604800000参数说明主队列 retention 设 1 天保证新数据优先重试队列 retention 设 7 天给异常数据足够的重试窗口。分区数主队列 32、重试队列 16按流量比例分配。副本数都是 3保证可用性。注意双队列的关键是「Worker 对 Queue1 的消费进度不被 Queue2 影响」。如果 RetryWorker 处理太慢导致 Queue2 积压不会反过来阻塞主队列。3. 实时性与可用性Storm 的 at least once 和降级开关怎么配合3.1 为什么选 at least once 而不是 exactly onceStorm 支持三种消息保证策略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) ) ENGINEInnoDB; -- 写入时用 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.workers8 \ -c topology.acker.executors8 # 运行时动态调整部分 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 切换到 s1s0 开放读写权限整个过程利用 MySQL 复制分发特性避免人工同步几分钟完成。结合 DB 降级功能只在 DNS 切换的几秒钟产生异常。-- 分片路由示意按 user_id 取模路由到不同库 -- 分片数 shardCount 42 的 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 一直不写 MySQLretry 队列持续积压。解决降级开关加自动超时机制比如打开后 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.resetlatest重启后从最新位置消费。如果需要精确恢复用 Kafka 的 offset 管理 API 手动提交和恢复。5.5 Redis 降级时吞吐量下降MySQL 压力过大现象Redis 降级后查询请求全部打到 MySQLMySQL 压力飙升响应变慢。原因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(hostredis-host, port6379) conn pymysql.connect(hostmysql-host, useruser, passwordpass, dbbehavior) def check_consistency(user_id): # 从 Redis 拉取行为记录 redis_data r.zrange(fbehavior:{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(fUser {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 生产者灌数据观察端到端延迟变化。注意压测时不要只打查询服务要同时灌数据。只打查询不灌数据测不出处理流的瓶颈。只灌数据不打查询测不出输出流的瓶颈。从那以后我每次设计实时链路都会强制走一遍「端到端延迟 数据一致性 压力测试」这三步。少一步上线后都可能出玄学问题。希望帮到你。本文还有配套的精品资源点击获取
返回列表