恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
携程实时用户行为服务架构:从埋点到决策的秒级链路设计
首页
资讯中心
/
携程实时用户行为服务架构:从埋点到决策的秒级链路设计
携程实时用户行为服务架构:从埋点到决策的秒级链路设计
发布时间:2026/10/3 20:17:51
简介这份文档面向后端架构师、大数据开发工程师及对实时计算感兴趣的技术人员系统梳理了携程实时用户行为服务系统的架构演进与实践经验。内容围绕推荐系统、动态广告、用户画像、浏览历史等场景剖析原有系统在数据覆盖、输出格式、日志处理性能上的痛点并给出重新设计后的处理流与输出流方案。文档详细阐述了Java、Kafka、Storm、Redis、MySQL、Tomcat、Spring技术栈的选型依据并从实时性、可用性、功能性、扩展性四个维度展开设计思路涵盖突发流量洪峰应对、双队列补偿重试、积压数据消解、DB降级等具体策略。资源为1个docx文件压缩包约161KB结构紧凑适合快速通读。目前已有130人学习读者可从中获取一套可借鉴的实时行为服务架构方法论与落地细节用于自身系统的设计参考与排错思路启发。1. 携程实时用户行为服务从埋点到决策的秒级链路用户在携程 App 上滑动酒店列表、点击房型、比价、下单这一串动作在 300 毫秒内就会被后端感知、聚合、打上标签然后推给推荐、风控、营销三个下游系统。这套链路就是携程实时用户行为服务系统架构要解决的问题把散落在客户端和服务端的埋点事件变成可查询、可订阅、可实时计算的用户行为流。它适合两类人看一类是正在做实时数仓或用户画像的工程师另一类是业务侧想接实时信号但不知道链路怎么搭的技术负责人。核心难点不在“收数据”而在“收得准、算得快、查得到、不丢不重”。下面按我实际落地过的路径从架构选型一路讲到参数调优和踩坑记录。2. 实时行为服务的分层架构与选型逻辑2.1 为什么不用纯 Kafka Flink 一把梭很多团队第一反应是 Kafka 收埋点、Flink 做计算、Redis 存结果三件套搞定。携程这种量级下这个方案会在两个地方翻车一是埋点事件类型超过 200 种每种事件的 schema 不同Kafka topic 如果按事件类型拆topic 数量爆炸运维成本极高二是下游消费方需求差异大推荐要的是用户最近 N 条行为序列风控要的是聚合计数营销要的是标签更新用同一个 Flink 作业输出所有形态状态会大到无法管理。常见做法是加一层“行为网关”做协议适配和路由。网关接收客户端 SDK 上报的原始事件做字段校验、事件归一化、用户 ID 映射device_id 到 user_id然后按“用户维度”而不是“事件类型维度”写入消息队列。这样 topic 数量可控下游按用户 ID 做 key 消费天然支持按用户聚合。选型上消息队列用 Kafka 还是 Pulsar我的判断标准是如果团队已经有成熟的 Kafka 运维体系不要为了 Pulsar 的分层存储特性去换迁移成本远大于收益。计算层 Flink 是默认选项但要注意 Flink 的状态后端选 RocksDB 还是内存直接决定你能扛多大数据量。存储层 Redis 做实时查询、HBase 做明细回溯、ClickHouse 做聚合分析三者不是替代关系是互补关系。2.2 行为网关的四个核心模块与配置行为网关是整条链路的入口它的稳定性决定一切。我一般把它拆成四个模块接入层、校验层、归一化层、路由层。接入层用 Netty 做 HTTP 和 WebSocket 双协议支持HTTP 用于批量上报WebSocket 用于实时性要求高的场景。线程模型上bossGroup 设 1workerGroup 设 CPU 核数 × 2不要设太大否则上下文切换开销会吃掉吞吐。// Netty 服务端启动配置 EventLoopGroup bossGroup new NioEventLoopGroup(1); EventLoopGroup workerGroup new NioEventLoopGroup(Runtime.getRuntime().availableProcessors() * 2); ServerBootstrap b new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) // 连接队列长度 .childOption(ChannelOption.TCP_NODELAY, true) // 禁用 Nagle 算法降低延迟 .childOption(ChannelOption.SO_KEEPALIVE, true) // 长连接保活 .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline().addLast(new HttpServerCodec()); ch.pipeline().addLast(new HttpObjectAggregator(65536)); // 聚合最大 64KB ch.pipeline().addLast(new BehaviorGatewayHandler()); } });这段代码的关键参数是SO_BACKLOG和HttpObjectAggregator的 maxContentLength。SO_BACKLOG设 1024 是经验值太小会导致高并发时连接被拒太大浪费内存。HttpObjectAggregator的 65536 字节限制意味着单次上报的 JSON body 不能超过 64KB超过会返回 413。如果你的埋点批量上报可能超过这个值要么调大要么在客户端做分片。校验层做两件事字段完整性检查和数据类型校验。我一般用 JSON Schema 定义每种事件的字段规范校验失败的事件不丢弃而是写入“死信队列”供后续排查。这里有个血泪经验不要在校验层做业务逻辑判断比如“用户必须登录才能上报”这种判断放到归一化层否则校验层会变得极其臃肿。归一化层负责把不同来源的事件统一成标准格式。比如 iOS 的page_id和 Android 的page_name要映射到同一个page_code时间戳统一成毫秒级 Unix 时间。这一步的配置建议用外部配置中心管理映射关系不要硬编码在代码里否则每次新增页面都要发版。路由层根据事件类型和用户 ID 决定写入哪个 Kafka topic。我一般按用户 ID 的 hash 取模分 64 个分区保证同一用户的事件有序。分区数不是越多越好64 个分区在 3 节点 Kafka 集群上每个节点约 21 个分区是吞吐和延迟的平衡点。2.3 Flink 作业的状态管理与 Exactly-Once 配置Flink 作业是计算核心它的状态管理直接决定系统能否做到 Exactly-Once。携程这种场景下我建议用 RocksDB 作为状态后端虽然读写性能比内存慢但支持的状态量级大得多。内存状态后端在用户行为序列超过 10 万条时就会 OOM。// Flink 环境配置与 Exactly-Once 开启 StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(60000); // 每 60 秒做一次 checkpoint env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); // 两次 checkpoint 最小间隔 env.getCheckpointConfig().setCheckpointTimeout(120000); // checkpoint 超时时间 env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 不允许并发 checkpoint env.setStateBackend(new RocksDBStateBackend(hdfs://namenode:8020/flink/checkpoints, true));enableCheckpointing(60000)这个 60 秒不是随便定的。太短会导致频繁写 HDFS影响吞吐太长则故障恢复时重放的数据量大恢复慢。我一般从 60 秒起步根据作业的背压情况调整。setMinPauseBetweenCheckpoints(30000)保证 checkpoint 之间至少间隔 30 秒避免上一个还没做完下一个就开始了。setMaxConcurrentCheckpoints(1)强制串行虽然会降低 checkpoint 频率但能避免状态不一致。RocksDB 的增量 checkpoint 要开启第二个参数true就是启用增量。全量 checkpoint 在状态达到 TB 级时每次都要几分钟增量能把时间压到秒级。但增量 checkpoint 的恢复依赖之前的 checkpoint 文件HDFS 上的文件不能随便清理要设置合理的保留策略。2.4 存储层选型Redis、HBase、ClickHouse 的边界存储层最容易犯的错是“一个 Redis 打天下”。Redis 适合存用户最近 N 条行为序列用 List 或 Sorted Set 结构查询延迟在毫秒级。但 Redis 的内存成本高存全量明细不现实。我一般只存最近 100 条行为更早的落到 HBase。HBase 的 RowKey 设计是核心。我一般用user_id reverse_timestamp作为 RowKey这样同一用户的最新行为排在前面查询时用Scan加setMaxResultSize就能快速拿到最近记录。reverse_timestamp 是Long.MAX_VALUE - timestamp保证时间倒序。列族只设一个列限定符用事件类型这样不同事件可以灵活扩展。ClickHouse 用于聚合分析比如“过去 7 天每个城市的酒店详情页 PV”。建表时用ReplacingMergeTree引擎按city_id和event_date分区排序键用(city_id, event_date, event_type)。查询时注意FINAL关键字的使用它会强制合并重复数据但性能开销大只在必要时用。提示Redis 和 HBase 的数据一致性靠 Flink 作业保证不要试图在存储层做双写同步那样会引入分布式事务的复杂度。3. 从埋点上报到实时查询的完整落地步骤3.1 客户端 SDK 的批量上报与重试策略客户端 SDK 是数据源头它的上报策略直接影响服务端压力。我一般配置为每 10 条事件或每 5 秒触发一次批量上报两者满足其一就发送。重试策略用指数退避首次失败后 1 秒重试第二次 2 秒第三次 4 秒最多重试 3 次。超过 3 次写入本地磁盘队列等下次启动时补报。// 客户端批量上报与指数退避重试 class BehaviorReporter { constructor(endpoint, batchSize 10, flushInterval 5000) { this.endpoint endpoint; this.batchSize batchSize; this.flushInterval flushInterval; this.queue []; this.retryCount 0; setInterval(() this.flush(), this.flushInterval); } report(event) { this.queue.push({ ...event, client_ts: Date.now() }); if (this.queue.length this.batchSize) this.flush(); } async flush() { if (this.queue.length 0) return; const batch this.queue.splice(0, this.batchSize); try { await fetch(this.endpoint, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify(batch) }); this.retryCount 0; } catch (e) { this.retryCount; if (this.retryCount 3) { const delay Math.pow(2, this.retryCount - 1) * 1000; setTimeout(() this.flush(), delay); } else { this.persistToLocal(batch); // 写入本地存储下次启动补报 } } } }batchSize设 10 是平衡点太小则请求数多太大则单次失败影响面大。flushInterval设 5000 毫秒保证即使事件少也能及时上报。指数退避的Math.pow(2, retryCount - 1) * 1000产生 1 秒、2 秒、4 秒的间隔避免雪崩。persistToLocal用 IndexedDB 或 localStorage注意容量限制一般只保留最近 1000 条。3.2 服务端 Kafka 生产者的 acks 与幂等配置服务端收到批量事件后归一化处理再写入 Kafka。生产者的配置直接决定数据不丢不重。Properties props new Properties(); props.put(bootstrap.servers, kafka1:9092,kafka2:9092,kafka3:9092); props.put(acks, all); // 所有 ISR 确认 props.put(retries, Integer.MAX_VALUE); // 无限重试 props.put(enable.idempotence, true); // 开启幂等 props.put(max.in.flight.requests.per.connection, 5); // 幂等要求 ≤5 props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.ByteArraySerializer); props.put(compression.type, lz4); // 压缩算法 props.put(linger.ms, 10); // 等待 10ms 凑批 props.put(batch.size, 16384); // 16KB 批大小acksall保证消息被所有同步副本确认配合enable.idempotencetrue实现生产者级别的 Exactly-Once。max.in.flight.requests.per.connection必须 ≤5否则幂等性失效。linger.ms10和batch.size16384是吞吐和延迟的折中如果业务对延迟极敏感可以把linger.ms设为 0但吞吐会下降。3.3 Flink 消费 Kafka 并写入 Redis/HBase 的代码骨架Flink 作业从 Kafka 消费做窗口聚合和序列拼接然后写入 Redis 和 HBase。DataStreamBehaviorEvent stream env.addSource( new FlinkKafkaConsumer(behavior-topic, new BehaviorDeserializer(), props) .setStartFromGroupOffsets() // 从消费者组偏移量开始 .setCommitOffsetsOnCheckpoints(true) // checkpoint 时提交偏移量 ); // 按用户 ID 分组维护最近 100 条行为序列 stream.keyBy(BehaviorEvent::getUserId) .process(new KeyedProcessFunctionString, BehaviorEvent, UserBehaviorSeq() { private ListStateBehaviorEvent state; Override public void open(Configuration parameters) { ListStateDescriptorBehaviorEvent descriptor new ListStateDescriptor(recent-behaviors, BehaviorEvent.class); state getRuntimeContext().getListState(descriptor); } Override public void processElement(BehaviorEvent event, Context ctx, CollectorUserBehaviorSeq out) throws Exception { state.add(event); ListBehaviorEvent list new ArrayList(); state.get().forEach(list::add); if (list.size() 100) { list list.subList(list.size() - 100, list.size()); state.update(list); } out.collect(new UserBehaviorSeq(event.getUserId(), list)); } }) .addSink(new RedisHBaseSink()); // 自定义 Sink同时写 Redis 和 HBasesetCommitOffsetsOnCheckpoints(true)保证偏移量和状态一起提交实现端到端的 Exactly-Once。ListState存最近 100 条行为超过就截断。RedisHBaseSink里先写 Redis 再写 HBaseRedis 写失败不影响 HBase但 HBase 写失败要抛异常触发重试。注意 Sink 的幂等性Redis 用LPUSHLTRIM保证只保留 100 条HBase 用Put天然幂等。3.4 实时查询接口的 Redis 数据结构设计下游查询接口从 Redis 读用户行为序列延迟要求 P99 在 10 毫秒以内。# Redis 存储结构示例 # 用户行为序列用 Listkey 为 user:behavior:{user_id} LPUSH user:behavior:12345 {event:hotel_detail,ts:1700000000000,page:hotel_123} LTRIM user:behavior:12345 0 99 # 只保留最近 100 条 # 用户标签用 Hashkey 为 user:tag:{user_id} HSET user:tag:12345 last_city 上海 last_active_ts 1700000000000 order_count_7d 3 # 聚合计数用 String INCRkey 为 metric:{event}:{date} INCR metric:hotel_detail:20240101List 的LPUSHLTRIM组合保证序列长度可控。Hash 存标签字段可动态扩展。String INCR 做计数注意设置过期时间比如EXPIRE metric:hotel_detail:20240101 86400 * 30保留 30 天。查询时用LRANGE user:behavior:12345 0 9拿最近 10 条不要一次拿 100 条减少网络传输。注意Redis 的 List 在元素过多时LTRIM会有性能开销建议在 Flink Sink 里做截断而不是依赖 Redis 的LTRIM。4. 避坑与排查实时行为服务最常见的五个翻车点4.1 现象Flink 作业频繁重启日志显示 Checkpoint 超时原因RocksDB 状态太大每次 checkpoint 要写 HDFS 的数据量超过checkpointTimeout设置。或者 HDFS 写入带宽不足导致 checkpoint 长时间完不成。解决先看 checkpoint 大小如果超过 10GB开启增量 checkpoint。如果已经开启检查 HDFS 的dfs.datanode.max.transfer.threads是否够大。还可以调大checkpointTimeout到 3000005 分钟但这是治标治本是减少状态量比如把用户行为序列从 100 条降到 50 条。4.2 现象Kafka 消费者组 Lag 持续增长但 Flink 作业没有背压原因Kafka 分区数和 Flink 并行度不匹配。比如 64 个分区Flink 并行度只有 8每个 subtask 要消费 8 个分区消费能力跟不上生产速度。解决把 Flink 并行度调到和分区数一致或者至少是分区数的约数。但并行度不能超过 TaskManager 的 slot 总数否则作业起不来。我一般设并行度 min(分区数, TaskManager 数 × 每 TM slot 数)。4.3 现象Redis 内存持续增长最终 OOM原因用户行为序列的 key 没有设置过期时间或者LTRIM没有生效。也可能是标签 Hash 的字段无限增加比如把每次请求的 trace_id 都写进去。解决所有 key 必须设 TTL行为序列设 7 天标签设 30 天。LTRIM要在每次LPUSH后立即执行不要攒一批再执行。标签字段要预定义不要动态加未知字段。4.4 现象HBase 写入变慢RegionServer 频繁 GC原因RowKey 设计不合理导致热点写入。比如用timestamp作为 RowKey 前缀所有写入集中在同一个 Region。解决RowKey 用hash(user_id) % 16 user_id reverse_timestamp把写入分散到 16 个 Region。reverse_timestamp保证查询时最新数据在前。另外HBase 的MemStore大小要调大默认 128MB 在高写入场景下不够我一般设 256MB。4.5 现象客户端上报的数据在服务端查不到但 Kafka 有消息原因归一化层把某些事件过滤掉了比如event_type不在白名单里。或者用户 ID 映射失败device_id找不到对应的user_id事件被丢弃。解决在归一化层加日志记录被过滤的事件和原因。用户 ID 映射失败的事件不要直接丢写入“待映射队列”等用户登录后补映射。白名单要可配置不要硬编码。5. 进阶技巧用 Flink SQL 做行为序列的实时特征计算Flink SQL 能大幅降低实时特征计算的开发成本。比如要算“用户过去 10 分钟内的酒店详情页点击次数”用 DataStream API 要写窗口、状态、定时器用 SQL 几行就搞定。-- 创建 Kafka 源表 CREATE TABLE behavior_events ( user_id STRING, event_type STRING, page_code STRING, event_ts BIGINT, WATERMARK FOR event_ts AS event_ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic behavior-topic, properties.bootstrap.servers kafka1:9092, format json, scan.startup.mode group-offsets ); -- 10 分钟滚动窗口计算每个用户的酒店详情页点击次数 SELECT user_id, COUNT(*) AS hotel_detail_pv, TUMBLE_START(event_ts, INTERVAL 10 MINUTE) AS window_start FROM behavior_events WHERE event_type hotel_detail GROUP BY user_id, TUMBLE(event_ts, INTERVAL 10 MINUTE);WATERMARK FOR event_ts AS event_ts - INTERVAL 5 SECOND定义水位线允许 5 秒乱序。TUMBLE是滚动窗口10 分钟一个窗口不重叠。scan.startup.mode group-offsets从消费者组偏移量开始保证重启后不丢数据。Flink SQL 的坑在于状态管理不如 DataStream 灵活。比如要维护用户最近 100 条行为序列SQL 的ARRAY_AGG不支持限制长度会无限增长。这种场景还是得用 DataStream API。我的习惯是聚合计数、简单窗口用 SQL序列拼接、复杂状态用 DataStream。另一个技巧是用 Flink SQL 的TEMPORAL JOIN做维表关联。比如行为事件要关联用户画像标签可以用FOR SYSTEM_TIME AS OF语法SELECT b.user_id, b.event_type, u.user_level, u.last_city FROM behavior_events b LEFT JOIN user_profile FOR SYSTEM_TIME AS OF b.event_ts AS u ON b.user_id u.user_id;FOR SYSTEM_TIME AS OF b.event_ts表示用事件时间对应的维表快照保证关联的是当时最新的标签而不是作业启动时的旧数据。维表用 JDBC 连接 MySQL 或 HBase注意设置lookup.cache.max-rows和lookup.cache.ttl控制缓存避免每次关联都查库。验证 Flink SQL 作业是否正确我一般用“双跑对比”同一份数据用 SQL 和 DataStream 各跑一遍对比结果。如果一致说明 SQL 逻辑正确。不一致就查窗口边界和水位线定义十有八九是水位线设得太激进导致迟到数据被丢。最后说一个我踩过的坑Flink SQL 的GROUP BY窗口聚合如果某个 key 在窗口内没有数据不会输出结果。但下游可能期望每个 key 每个窗口都有一条记录哪怕是 0。解决方法是把窗口聚合结果再LEFT JOIN一个全量 key 表用COALESCE把 NULL 转成 0。这个操作在 DataStream 里用sideOutput也能做但 SQL 更直观。希望帮到你。本文还有配套的精品资源点击获取