
简介这份文档面向后端架构师、大数据开发工程师及对实时计算感兴趣的技术人员围绕携程实时用户行为服务系统的架构实践展开重点解决原有系统数据覆盖不全、输出格式不统一、日志处理模块难以扩容等痛点。资源为单个docx文件压缩包约161KB内容涵盖架构设计背景、处理流与输出流双流向设计、JavaKafkaStormRedisMySQL技术栈选型以及实时性、可用性、功能性与扩展性四个维度的具体方案。读者可从中了解双队列补偿重试、积压数据消解、DB降级等实战思路并获取每天处理约20亿数据、上线可用约300毫秒、查询平均延迟约6毫秒等性能指标参考。目前已有130人学习适合需要构建或优化实时用户行为服务的技术团队借鉴。1. 携程实时用户行为服务系统架构实践从埋点到决策的秒级闭环用户在携程 App 里滑动列表、点开酒店详情、反复比价、加购又放弃这些动作每时每刻都在产生数据。问题在于这些行为数据如果只是躺在离线数仓里等 T1 跑批那它对搜索排序、推荐、风控、实时营销的价值就废了一大半——等你知道用户刚才看了三次某家酒店用户早就下单走人了。携程实时用户行为服务系统架构实践要解决的就是把这堆高频、无序、量级巨大的行为事件在秒级内变成可查询、可订阅、可驱动业务决策的服务能力。这套系统适合谁看如果你是正在做用户行为采集、实时数仓、特征平台、推荐系统在线特征供给的工程师或者你所在团队正准备把离线用户画像升级成实时版本那这套架构思路可以直接对照落地。它不追求某个单点技术的极致性能而是强调在携程这种多业务线、多端、高并发的场景下如何把采集、传输、计算、存储、服务五层串成一条稳定链路。下面我会按「先讲清楚每层为什么这么选再给可复现的配置和代码最后把踩过的坑摊开」的节奏展开中间会涉及 Kafka、Flink、HBase、Redis 这些常见组件的组合方式也会给出参数建议和排查手段。2. 行为采集与传输层SDK 埋点规范与 Kafka 管道设计2.1 为什么采集层要统一 SDK 而不是各业务自己上报携程的业务线多酒店、机票、火车票、门票各自有前端团队。如果每个团队自己定义埋点格式、自己选传输通道最后到实时计算层就是一场灾难字段名不统一、时间戳精度不一致、事件类型枚举冲突。所以采集层的第一原则是统一 SDK 统一事件模型。常见做法是提供一个跨端 SDKiOS/Android/Web/小程序内部封装三件事事件定义注册、本地缓存与批量发送、失败重试与降级。事件模型至少包含以下字段字段名类型说明event_idstring全局唯一事件 ID用于去重event_typestring事件类型如 page_view、click、order_submituser_idstring登录用户 ID未登录时为空device_idstring设备指纹未登录时作为主键session_idstring会话 ID由 SDK 生成page_idstring页面标识element_idstring被点击元素标识ts_clientbigint客户端发生时间毫秒ts_serverbigint服务端接收时间毫秒propertiesmap业务自定义属性这个模型的关键在于user_id 和 device_id 必须同时存在后续实时计算才能做匿名到登录的身份归一。ts_client 和 ts_server 的差值可以用来监控端到端延迟也能识别客户端时间被篡改的情况。2.2 Kafka 主题划分与分区策略采集 SDK 把事件批量发到接入网关网关做基础校验后写入 Kafka。Kafka 的主题设计直接决定后续 Flink 作业的并行度和消费顺序。我一般会按「业务域 事件大类」拆主题而不是所有事件塞一个 topic。比如behavior_page_view页面浏览事件behavior_click点击事件behavior_order订单相关事件behavior_search搜索行为事件分区数怎么定一个经验公式分区数 峰值 QPS / 单分区消费能力。单分区消费能力取决于 Flink 算子复杂度纯过滤场景单分区能到 5000 QPS 左右带状态计算会降到 1000 到 2000。携程大促期间行为事件峰值能到百万 QPS 级别所以分区数通常按业务域给到 64 到 128 个。分区键的选择有个坑如果用 user_id 做分区键未登录用户 user_id 为空会导致大量数据涌向同一个分区。正确做法是用user_id device_id拼接后的哈希值做分区键保证同一用户的事件落到同一分区同时未登录用户也能均匀分布。// 分区键生成逻辑保证同一用户事件有序且分布均匀 public String buildPartitionKey(String userId, String deviceId) { // 登录用户优先用 userId未登录用 deviceId String key (userId ! null !userId.isEmpty()) ? userId : deviceId; // 加盐避免热点用户导致单分区过载 int salt Math.abs(key.hashCode()) % 16; return salt _ key; }这段代码的逻辑是先用 userId 或 deviceId 确定用户标识再通过取模加盐把原本可能集中在少数分区的热点用户打散。参数 16 是盐桶数量可以根据分区数调整一般取分区数的 1/8 到 1/4。加盐后同一用户的事件仍然落在同一分区因为盐值由 key 本身决定是确定性的。2.3 接入网关的限流与降级网关层必须做限流否则某个业务线发版时埋点暴增会把整个 Kafka 打满。我一般用令牌桶做单业务线限流桶容量按该业务线历史峰值的 1.5 倍设置令牌生成速率按历史均值的 2 倍设置。超过限流阈值的事件直接写本地磁盘做降级等流量回落再补发。注意降级磁盘文件要设置过期时间比如 24 小时避免磁盘被写满。补发时要带上原始 ts_client否则实时计算层的时间窗口会错乱。3. 实时计算层Flink 作业的状态管理与窗口设计3.1 为什么选 Flink 而不是 Spark Streaming在携程这个场景下行为事件需要做会话窗口聚合、CEP 模式匹配比如「搜索→点击→加购→未支付」的漏斗、以及维表关联关联用户画像、酒店基础信息。Spark Streaming 的微批模型在会话窗口和 CEP 上表达力弱而且状态管理不如 Flink 原生。Flink 的 EventTime 语义和 Watermark 机制能较好地处理客户端时间乱序的问题这对行为数据尤其重要——移动端网络抖动会导致事件到达顺序和发生顺序不一致。选型确定后作业的并行度设置是第一个要调的参数。并行度一般设为 Kafka 分区数的整数倍保证每个分区都有消费者。比如 topic 有 128 个分区Flink 并行度可以设 128 或 256。设太小会导致消费积压设太大则增加状态后端压力。3.2 会话窗口与去重逻辑用户行为分析里最常用的窗口是会话窗口。携程场景下会话超时时间一般设 30 分钟——用户 30 分钟没有新事件就认为会话结束。Flink 的EventTimeSessionWindows.withGap(Time.minutes(30))可以直接表达。但会话窗口有个问题如果用户一直有事件窗口永远不触发。所以实际生产中我会加一个最大会话时长限制比如 2 小时强制触发一次。代码结构如下DataStreamBehaviorEvent stream env .addSource(new FlinkKafkaConsumer(behavior_click, schema, props)) .assignTimestampsAndWatermarks( WatermarkStrategy.BehaviorEventforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, ts) - event.getTsClient()) ); DataStreamSessionAggregate sessionResult stream .keyBy(event - event.getUserId() ! null ? event.getUserId() : event.getDeviceId()) .window(EventTimeSessionWindows.withGap(Time.minutes(30))) .trigger(ContinuousEventTimeTrigger.of(Time.minutes(10))) // 每10分钟触发一次增量输出 .aggregate(new SessionAggregator(), new SessionWindowFunction());这段代码的关键参数有三个forBoundedOutOfOrderness(Duration.ofSeconds(5))表示允许 5 秒乱序超过 5 秒的事件会被丢弃到侧输出流withGap(Time.minutes(30))是会话超时ContinuousEventTimeTrigger.of(Time.minutes(10))让窗口每 10 分钟输出一次增量结果避免长会话导致结果延迟。去重逻辑放在聚合函数里用event_id做去重键状态保留时间设为会话超时时间的两倍。Flink 的ValueState可以存最近处理过的 event_id 集合但要注意状态大小——如果单个会话事件量很大可以用布隆过滤器替代精确集合。3.3 维表关联的异步 IO 优化行为事件需要关联用户画像和酒店信息这些维表存在 HBase 或 MySQL 里。如果用同步查询每条事件都阻塞等待吞吐量上不去。Flink 的 Async I/O 可以解决这个问题但要注意三点第一异步客户端要复用连接池不能每次查询新建连接。第二超时时间要设合理一般 100 到 200 毫秒超时后走降级逻辑返回空维度。第三要监控异步队列的积压情况如果队列满了说明下游维表查询能力不足需要扩容维表存储或加缓存。// 异步维表关联示例 AsyncDataStream.unorderedWait( stream, new AsyncUserProfileLookup(), 200, TimeUnit.MILLISECONDS, // 超时时间 1000 // 最大异步请求数 ).setParallelism(128);参数1000是单个并行度下允许的最大异步请求数乘以并行度 128 就是总并发查询量。这个值要根据维表存储的承载能力调整HBase 集群一般能扛住几万 QPSRedis 能到十万级别。4. 存储与服务层HBase 行键设计与 Redis 热点缓存4.1 HBase 行键怎么设计才能支撑实时查询实时计算的结果需要落地存储供业务方查询。携程场景下查询模式主要有两种按用户查最近行为序列按酒店查实时热度。HBase 适合存用户行为明细行键设计是核心。常见错误是用user_id 时间戳做行键这样会导致同一用户的数据集中在一个 Region用户量大时出现热点。正确做法是加盐或反转。我一般用salt user_id ts作为行键salt 取 user_id 哈希后对 16 取模。这样同一用户的数据分散在 16 个 Region 上查询时用scan加startRow和stopRow限定范围。列族设计上行为明细用一个列族info列限定符用事件类型版本数设为 1 或 3。版本数设 1 表示只保留最新值设 3 表示保留最近三个版本。行为序列查询通常需要多个事件所以更常见的做法是每行存一个事件用时间戳做排序。4.2 Redis 缓存热点用户行为HBase 查询延迟在 10 到 50 毫秒对于推荐系统在线特征获取来说还是偏慢。所以会在 Redis 里缓存最近活跃用户的行为序列。缓存结构用Sorted Setscore 是事件时间戳member 是事件 ID 或序列化后的事件内容。查询时用ZREVRANGE取最近 N 条。缓存更新策略是「写时更新 读时回填」。Flink 作业计算出会话结果后除了写 HBase还异步写 Redis。读端如果 Redis 未命中从 HBase 查完再回填 Redis设置 10 分钟过期时间。# Redis 写入逻辑使用 pipeline 批量提交 import redis import json r redis.Redis(connection_poolpool) pipe r.pipeline() for event in session_events: key fuser:behavior:{event[user_id]} score event[ts_client] member json.dumps(event, ensure_asciiFalse) pipe.zadd(key, {member: score}) pipe.expire(key, 600) # 10分钟过期 pipe.execute()这段代码用 pipeline 减少网络往返zadd的 score 用客户端时间戳保证排序正确expire设置 600 秒过期避免冷用户占用内存。注意 member 用 JSON 序列化后可能较大如果事件属性多可以只存关键字段。4.3 服务接口的降级与熔断对外提供的查询接口要有降级策略。当 HBase 或 Redis 响应超时接口应该返回最近一次缓存的结果或空结果而不是直接报错。熔断用 Hystrix 或 Sentinel 实现阈值一般设 50% 错误率或 500 毫秒平均响应时间熔断后 30 秒半开重试。5. 避坑与排查实时行为系统最常见的五个翻车现场5.1 现象Flink 作业频繁反压Kafka 消费积压原因通常是某个算子处理能力不足比如维表关联的异步请求超时导致队列满或者窗口状态过大导致 checkpoint 时间过长。排查时先看 Flink Web UI 的反压指标定位到具体算子再看该算子的 CPU 和内存使用率。解决如果是异步 IO 问题调大超时时间或增加维表存储容量如果是状态问题检查状态后端配置把 RocksDB 的write_buffer_size和max_background_jobs调大或者缩短状态保留时间。5.2 现象用户行为序列查询结果乱序原因HBase 行键设计时时间戳精度不够或者 Redis Sorted Set 的 score 用了服务端时间而非客户端时间。客户端时间乱序会导致排序错误。解决统一用客户端毫秒时间戳做排序键HBase 行键里时间戳倒序存储用Long.MAX_VALUE - tsRedis score 直接用 ts_client。5.3 现象大促期间 Kafka 写入延迟飙升原因分区键设计不合理导致热点分区或者网关限流阈值设太高瞬时流量打满 Kafka 磁盘 IO。解决检查分区键分布用kafka-topics.sh --describe看各分区 leader 的写入速率。如果确实热点增加盐桶数量或改用复合键。限流阈值按历史峰值 1.5 倍设置超过部分走降级磁盘。5.4 现象Flink checkpoint 失败作业重启后状态丢失原因状态后端存储HDFS 或 S3写入超时或者 checkpoint 间隔太短导致频繁触发。携程场景下状态可能到 TB 级别checkpoint 时间会很长。解决checkpoint 间隔设为 5 到 10 分钟超时时间设为 10 分钟。开启增量 checkpointRocksDB 支持减少每次上传的数据量。同时监控 checkpoint 的 duration 和 size如果持续增长说明状态泄漏要检查是否有未清理的过期状态。5.5 现象Redis 内存暴涨缓存命中率下降原因过期时间设置不合理或者 key 数量太多导致内存碎片。行为序列的 member 如果存了完整事件 JSON单个 key 可能几十 KB百万用户就是几十 GB。解决只存关键字段或者用压缩算法如 snappy压缩后再存。设置 maxmemory-policy 为 allkeys-lru让 Redis 自动淘汰冷数据。监控used_memory和evicted_keys如果淘汰量持续大于 0说明内存不足需要扩容或缩短过期时间。6. 进阶技巧用侧输出流做延迟数据处理与实时监控埋点实时行为系统里延迟数据是绕不开的问题。移动端网络抖动、SDK 重试、网关降级补发都会导致事件到达时间远晚于发生时间。如果直接丢弃会丢失部分用户行为如果无限等待又会影响窗口触发。我一般用 Flink 的侧输出流Side Output把超过 Watermark 的事件收集起来单独走一个延迟处理通道。具体做法是定义一个OutputTagBehaviorEvent在ProcessFunction里判断事件时间是否小于当前 Watermark 减去允许延迟如果是则ctx.output(tag, event)输出到侧流。侧流可以再接一个窗口窗口大小设为 1 小时处理那些迟到但仍有价值的事件。这些事件通常用于离线校正不参与实时决策但可以更新用户画像的长期特征。OutputTagBehaviorEvent lateTag new OutputTagBehaviorEvent(late-events){}; SingleOutputStreamOperatorSessionAggregate mainStream stream .keyBy(...) .window(...) .allowedLateness(Time.minutes(5)) // 允许5分钟迟到 .sideOutputLateData(lateTag) .aggregate(...); DataStreamBehaviorEvent lateStream mainStream.getSideOutput(lateTag); // 延迟数据写入单独的 Kafka topic供离线校正使用 lateStream.addSink(new FlinkKafkaProducer(behavior_late, schema, props));这段代码里allowedLateness(Time.minutes(5))表示窗口触发后仍保留状态 5 分钟期间到达的迟到数据会重新触发窗口计算。超过 5 分钟的进入侧输出流。参数 5 分钟是根据业务容忍度定的携程场景下大部分事件延迟在 3 秒内5 分钟已经能覆盖 99.9% 的情况。另一个进阶技巧是把监控埋点做进实时链路。每个算子输出一条监控事件到单独的 topic包含算子名称、处理条数、平均延迟、状态大小。这些监控事件再被一个独立的 Flink 作业消费写入时序数据库配合告警规则。这样不用等业务方反馈自己就能发现链路异常。我自己的习惯是每次上线新的 Flink 作业先跑 24 小时观察 checkpoint 和反压指标确认稳定后再逐步调大并行度。实时行为系统的稳定性比峰值性能更重要因为一旦链路断了推荐和风控都会受影响。希望帮到你。本文还有配套的精品资源点击获取