Flink实时标签双写StarRocks与PolarDB-X的实践与调优
发布时间:2026/9/14 20:08:09
分类:文化教育
浏览:1234

1. 场景与整体设计思路1.1 为什么需要双写而不是单写再拷贝做实时标签系统的时候我遇到的最麻烦的一件事不是标签逻辑本身而是怎么把同一份实时结果同时喂给两个用途完全不同的存储引擎。标签数据一旦生成在线服务要毫秒级查到最新值运营分析要能按人群、按标签维度做聚合统计。这两个需求放在同一个库里谁也扛不住——在线查询要的是低延迟和简单主键扫描分析统计要的是高吞吐导入和列存扫描能力。所以当时设计目标非常明确一份 Flink 实时作业从 Kafka 消费用户行为事件流经过标签计算逻辑后双写进入 StarRocks 和 PolarDB-X。StarRocks 负责分析侧供 BI 报表、人群圈选、运营看板使用PolarDB-X 负责在线侧承接 C 端查询、规则引擎、实时风控等对延迟和一致性要求很高的场景。这个方案的直接收益是省掉了一套“先写 A 库再同步到 B 库”的中间链路。早先的架构是 Flink 写一张宽表到 StarRocks然后靠 StarRocks 的外部表或者定时任务再同步到 PolarDB-X链路长、延迟高、失败点还多。双写看起来只是多加了一个 Sink实际是把数据分发的责任从存储层转移到了计算层让 Flink 统一管理两个目标的写入语义、幂等性和失败恢复。1.2 为什么选 StarRocks 和 PolarDB-X 这两类引擎选型这件事网上对比文章很多但大部分都是拿压测数据说话忽略了业务形态。StarRocks 在分析侧的优势是导入链路成熟支持 Stream Load、Broker Load、Routine Load而且主键模型可以做实时更新对标签这种高频变更的数据非常友好。PolarDB-X 则是对标 MySQL 生态的分布式关系库兼容 MySQL 协议业务代码几乎零改造就能接入而且支持分布式事务和全局二级索引适合在线交易类查询。这两个引擎的定位完全不同StarRocks 是 OLAPPolarDB-X 是 OLTP。双写作业的核心权衡也就在这里一份数据流要同时适配分析引擎的批量导入语义和在线引擎的单条 upsert 语义。Flink 的 JDBC 连接器天然适合后者而 StarRocks 官方提供的 Flink connector 走的是 Stream Load 通道两者在吞吐模型、事务机制、失败重试策略上完全不同不能简单套用一套参数。注意不要试图用统一的抽象封装两个 Sink 的写入逻辑。StarRocks 的 Stream Load 是按批次提交的PolarDB-X 的 JDBC upsert 是按缓冲行数 flush 的抽象层如果做得太厚反而掩盖了各自真正的调优入口。1.3 标签双写作业的整体数据流整个作业的拓扑其实不复杂Source 是 Kafka中间是标签计算逻辑Sink 是两个目标端。但真正落地时每一环都有值得抠的细节。Kafka 里存放的是用户行为事实数据比如浏览、点击、加购、下单事件。标签计算层用 Flink SQL 做 Segment 聚合和时间窗口统计产出的是用户维度的标签结果比如“近 7 天加购次数”“高活跃用户”“价格敏感人群”等。计算完成之后数据以用户 ID 为主键、标签字段为列的一行一行的形式分发给两个 Sink。这里有个容易忽略的点标签数据是业务状态不是事件日志。事件日志只需要 append标签数据则需要按主键覆盖更新。这就导致 Sink 层的写入语义必须是 upsert而不仅仅是 insert。Stream Load 端可以通过 StarRocks 主键模型实现PolarDB-X 端则依赖 JDBC 连接器的ON DUPLICATE KEY UPDATE能力。这两个机制一个走的是“批量导入 主键去重”一个走的是“逐条 upsert 事务批量提交”性能差异会在后面章节专门展开。2. 双写作业的核心链路拆解与参数选型2.1 Kafka Source 的并行度设置与消费语义Kafka Source 是整个作业的入口并行度设置直接影响两个 Sink 的写入压力。很多人直接把并行度拉到和分区数一样但实际上双写作业的瓶颈往往不在 Source而在 Sink 的写入吞吐。如果 Source 并行度太高Kafka 拉取速度远超存储端的写入能力checkpoint 会频繁超时作业最终被反压拖死。我当时的做法是先压测单并行度下 StarRocks Stream Load 的吞吐上限再反推 Source 并行度。假设单并行度 Stream Load 能跑到 30MB/s两个 Sink 同时写吞吐就是 15MB/s 左右Kafka 单分区消费速率大约 10MB/s那么并行度设为 46 个消费者基本能打满存储端又不至于过度拉数据。消费语义方面双写场景我强烈建议用commit_on_checkpointtrue也就是 checkpoint 完成才提交 Kafka offset。这样 Flink 的 Exactly-once 语义虽然不能保证两个外部存储端与 Kafka offset 的原子一致但至少能保证“重启后不丢数据”的最底线。配合两个 Sink 的幂等写入最终能做到“最多一次 幂等 不重不丢”的效果。Kafka 侧还需要留意消息体大小。标签中间结果经常带有 JSON 嵌套结构比如人群标签的命中规则明细单个消息可能到几十 KB。如果默认的fetch.max.bytes太小消费延迟会肉眼可见地升高。我建议把kafka.source.fetch.max.bytes调到 50MB、kafka.source.max.partition.fetch.bytes调到 10MB减少小包频繁拉取的网络开销。2.2 标签数据的 Schema 设计与主键策略双写作业最坑的一点是表结构设计。StarRocks 的主键模型和 PolarDB-X 的普通 MySQL 表结构虽然都是主键去重但对字段类型、索引方式的要求完全不同。StarRocks 主键模型要求主键字段尽量短不建议把 user_id 和标签维度混成一个超长联合主键否则主键索引占用内存会非常大。PolarDB-X 虽然对主键长度没那么敏感但分布式分区键的选择会直接影响写入热点的分布。我最终的设计方案是两个库的主键都是user_id tag_date。user_id 是 biginttag_date 是日期字符串这样既保证标签按天可回溯也避免主键无限膨胀。StarRocks 侧再把 tag_date 设为分区字段每天一个分区PolarDB-X 侧把 user_id 作为分区键保证同一个用户的标签查询落在同一个分片上。字段层面标签列我建议设计成“统一扩展字段 明细拆分字段”的组合。扩展字段是一个 map 类型或者 JSON 字符串承接临时新增的标签不用改表结构明细字段是固定的高频标签比如is_high_active、buy_tendency_score方便查询和索引。实操心得实时标签千万别一开始就追求把几十个标签全部拆成独立列。标签的迭代非常快今天新增一个“新品偏好标签”明天新增一个“大促敏感标签”频繁 ALTER TABLE 在大数据量下代价很高。先把高频核心标签独立成列其余塞进 JSON 扩展字段等业务验证某个标签确实值得独立列存储时再正式拆出。2.3 标签计算逻辑与 Flink SQL 的 Watermark 处理Flink SQL 处理标签计算最核心的是 Watermark 策略。标签场景和普通实时报表不一样报表晚到数据延迟几分钟展示影响不大但标签一旦基于错误时间窗口计算会影响后续所有的圈人和营销动作。我的 Watermark 策略是允许 30 秒乱序即WITH WATERMARK FOR event_time AS event_time - INTERVAL 30 SECOND。这背后有两个考虑Kafka 到 Flink 的链路延时正常情况下在 100ms 以内业务埋点时钟偏差一般在秒级30 秒是一个不会触发大量窗口重算又能容忍大部分乱序的安全值。如果业务里有离线导入的补数场景这个值就需要放大到 5 分钟以上但那样实时性会明显下降需要业务侧接受。标签计算里还有一个容易踩坑的点多流 Join 的维表补全。比如用户行为流需要 join 用户维表获取用户的注册城市、会员等级。维表数据是从 MySQL CDC 同步到 Kafka 的如果用 Flink SQL 的普通 temporal join维表会存在 TM 堆内存里量大一点直接 OOM。后来我改成了基于 RocksDB 的 lookup join状态后端调成 RocksDB才把 TM 内存稳定下来。这个改造会牺牲一部分 join 的吞吐但双写作业本来 Sink 侧就是瓶颈join 慢一点反而让后端压力更可控。3. 双写实现的工程要点两个 Sink 的不同性格3.1 StarRocks SinkStream Load 的参数调优StarRocks 的 Flink connector 底层是攒批后通过 HTTP 接口发 Stream Load 请求。它的吞吐跟两个因素强相关单次导入的数据量和并发导入的 Stream Load 任务数。先说单次导入的数据量。攒批不是越多越好Stream Load 请求有超时时间默认 30 秒如果单批超过几万行、几十 MB导入时长可能超过超时阈值导致任务报错重试。我当时把sink.properties.column_separator设成\t用于减少转义开销sink.buffer-flush.max-rows设为50000行sink.buffer-flush.max-bytes设为10MBsink.buffer-flush.interval设为5s。这套参数在单并行度下能稳定跑出 20MB/s 以上的写入速度。再说并发。StarRocks 虽然支持多个 Stream Load 同时写入一张表但并发数过高时 BE 节点 CPU 和磁盘 IO 会急剧飙升Label 冲突也可能导致重复导入报错。我建议 Stream Load 的并发控制在Tablet 数量 / 2以内不要盲目设置成和 Flink Sink 并行度一致。还有一个特别重要的点StarRocks 主键模型的写入性能对写入乱序非常敏感。如果同一个用户 ID 的标签在这批数据里出现在前面的分片下一批又出现在后面的分片StarRocks 内部要花大量资源做主键去重和版本合并。解决方法是把 Flink Sink 并行度设为与user_id的取模分区数一致保证同一用户的数据尽量落在同一个 Sink 子任务、同一批 Stream Load 里。3.2 PolarDB-X SinkJDBC upsert 的 Batch 策略PolarDB-X 兼容 MySQL 协议Flink 官方 JDBC 连接器可以直接用。但官方连接器的默认参数是为普通 MySQL 设计的直接用在分布式数据库上会出不少问题。核心参数是这几个sink.buffer-flush.max-rows默认 100 条对于分布式库来说太小了我调到 1000 条。sink.buffer-flush.interval默认 1 秒调大到 3 秒。sink.max-retries默认 3 次我调到 5 次。这里的关键是理解分布式事务在 Flink JDBC Sink 里的表现。Flink 官方 JDBC Sink 的 upsert 实现是攒批后走addBatch/executeBatch这个批在 PolarDB-X 上会变成一个分布式事务。批次太大事务执行时间变长冲突概率增大死锁回滚的概率也高批次太小网络开销占比大吞吐上不去。1000 条一批在user_id分区键设计合理的情况下一般事务执行时间在 100ms 左右是一个比较平衡的值。PolarDB-X 的 upsert 语法走 MySQL 语义即INSERT INTO table VALUES (...) ON DUPLICATE KEY UPDATE col VALUES(col)。Flink JDBC 连接器需要把连接参数里加上rewriteBatchedStatementstrue同时在 SQL 初始化语句里写清楚主键更新逻辑。如果不加rewriteBatchedStatementsMySQL JDBC 驱动不会真正批量执行而是逐条 prepareStatement性能相差 5 倍以上。3.3 双写一致性事务边界与幂等设计双写的核心痛点不是“写不进去”而是“写进去了但两边不一致”。由于两个 Sink 是独立的Flink checkpoint 成功只能保证两个 Sink 各自把攒批数据刷给了下游存储但 Kafka offset 提交、StarRocks 导入完成、PolarDB-X 事务提交这三件事不在同一个原子边界内。作业在写 PolarDB-X 成功后、写 StarRocks 失败前发生了故障就会出现一边有新数据、一边是旧数据的中间状态。解决这个问题的关键是幂等 对账不可能是真正的分布式原子事务。幂等层的设计如下StarRocks 侧Stream Load 使用label保证幂等Flink connector 的 label 默认是{table}_{uuid}我改成了基于 checkpoint_id 生成这样 Flink 恢复时会自动跳过已导入的批次。PolarDB-X 侧依靠主键 upsert 天然幂等同一行数据无论写几次最终结果一致不会产生重复记录。两个 Sink 之间不追求事务同步而是依赖 Kafka 的 offset 作为逻辑时间线通过定时任务对比两边的max(tag_date)和数据行数发现不一致。这样设计之后即使发生故障恢复流程也是“Flink 作业从最近 checkpoint 重启重放部分 Kafka 消息两个 Sink 各自动过滤已经写入的数据”最终两边收敛到一致状态。4. 性能权衡实录同一份数据两种写入模型4.1 吞吐与延迟的实测数据对比双写作业是否达标不能只看 Flink 作业的吞吐指标要分两条链路分别测。我当时用模拟数据压测的结果大致如下指标StarRocks 主键模型PolarDB-X 分区表写入模型Stream Load 批量导入JDBC batch upsert单并行度峰值吞吐25 MB/s 左右4000 行/s 左右单条数据平均写入延迟批次提交后秒级可见批次提交后毫秒级可见对 CPU 的消耗相对较低HTTP 传输 BE 合并相对较高事务解析 索引维护对内存的消耗Flink TM 端攒批内存可控JDBC 驱动内部缓冲需额外关照反压敏感度低Stream Load 失败会阻塞但可重试高事务超时或死锁会导致连续失败这张表其实很好地说明了为什么选型要按业务需求来。StarRocks 擅长的是大吞吐批量写入延迟在秒级适合分析场景PolarDB-X 的定位是近实时在线查询延迟在毫秒级但吞吐天花板远低于 StarRocks。4.2 双写带来的 CPU 放大效应双写最直接的代价是 Flink 作业的 CPU 占用比单写高了一大截。原因有两个第一同样的序列化、网络传输和内存拷贝要做两遍。Flink 内部对每个 Sink 都会做独立的 serializer 和 network buffer这就意味着数据从算子产出来后会同时被拷贝给两个下游。如果 Sink 并行度设置不一样还会引入数据倾斜和网络 shuffle。第二JDBC 连接器的ON DUPLICATE KEY UPDATE需要 Flink 端做字段值的反射拼接这个操作对 CPU 的消耗比预想中高很多尤其在字段数超过 20 个以后尤其明显。我当时的优化方案是PolarDB-X 需要 upsert 的字段从一开始就控制在 10 个以内其余字段能不进就不进。分析字段全部放 StarRocks在线查询只需要最核心的几个标签。4.3 反压控制策略用 checkpoint 指标反推并行度双写作业上线后我遇到最多的是反压问题而且反压源头基本都是 PolarDB-X 这个 JDBC Sink。因为 JDBC Sink 的攒批 flush 有网络往返和数据库事务执行时间一旦 PolarDB-X 侧出现慢 SQL 或锁等待flush 时长会从几十毫秒飙升到几秒甚至超时Sink 算子背压很快传遍全链路。排查反压不能只看 Kafka lag要结合 Flink UI 的BackPressure和Checkpoint指标一起看。当时我总结了一个判断流程如果Checkpoint Duration持续超过Checkpoint Interval说明 Sink flush 已经阻塞主线。打开 Flink Web UI 的 BackPressure 页面定位到具体是哪个 Sink 子任务处于 High 状态。看对应 TaskManager 的线程 dump确认是阻塞在 HTTP 的 Stream Load 请求上还是 JDBC 的 executeBatch 上。如果是 JDBC Sink先看数据库侧有没有锁等待和慢查询没有的话再调低sink.buffer-flush.max-rows并增大 Sink 并行度。这套排查流程帮我定位了多次并非 Flink 本身问题、而是 PolarDB-X 一个二级索引导致插入慢的系统性问题。注意反压不是只能靠调大并行度硬扛。先确认反压是“下游存储能力不够”还是“单批写入效率太低”两者的调优方向完全相反。前者要增加并行度后者要增大批次大小。5. 常见问题与排查技巧实录5.1 Flink JDBC 连接器异常Failed to upsert data的根源这个报错在双写作业刚上线时几乎天天出现。排查下来发现是两类原因第一类原因是 JDBC 连接器默认的rewriteBatchedStatementsfalse。PolarDB-X 作为分布式数据库驱动收到没有 rewrite 标记的 batch 请求时实际上是逐条执行每条数据一次网络往返。当批次数据量到几百上千条时执行时间会被网络延迟放大导致网关超时。解决办法很简单在 JDBC URL 加上rewriteBatchedStatementstrueuseServerPrepStmtsfalse。第二类原因是 SQL 初始化语句里的字段顺序与上游 DataStream 的 Row 类型字段顺序不一致。Flink JDBC 连接器按位置映射参数如果初始化 SQL 里写(user_id, tag_json)但上游 Row 的顺序是(tag_json, user_id)参数就串位了。这种错误在字段少的时候很难发现因为类型如果恰好都是 string数据就会写进错误的列里。建议在定义JdbcExecutionOptions之前先用简单的 Print Sink 把 Row 的字段顺序打印出来核对一遍。5.2 StarRocks 导入抖动Publish Timeout与 Label 重复StarRocks Stream Load 的Publish Timeout是我生产环境里遇到最多的一类错误。它的直接原因是 BE 节点执行事务 Publish 阶段超时常见诱因有三个某个 BE 节点磁盘 IO 过高导致版本合并变慢。单次导入的数据分布太散涉及太多 Tablet。并发 Stream Load 任务过多超过 BE 节点处理能力。Publish Timeout 的排查思路是先看 StarRocks FE 的审计日志找到对应 Label 的请求耗时是卡在LOAD阶段还是PUBLISH阶段。卡在LOAD说明是写入阶段慢需要检查磁盘卡在PUBLISH说明是版本合并慢需要降低并发或者缩小单批大小。关于 Label 重复的问题Flink connector 在任务重启后会用新的 UUID 作为 Label理论上不会重复。但如果发生了“task 提交后 FE 端标记成功、Flink 端没收到响应”的情况FE 端可能出现大量已经成功导入但未被确认的 Label长期积累会影响 FE 的内存。我建议定期清理 StarRocks 端超过保留时间的历史 Label或者用脚本清理已完成的导入事务。5.3 Kafka 消息延迟高的隐蔽原因双写作业链路里Kafka 消息延迟高不一定就是 Kafka 本身的问题。有段时间我们的监控面板上 Kafka Consumer Lag 持续在涨但 Flink 作业的 CPU、内存都正常Sink 侧也没有反压。后来查了半天才发现是 Source 端的空闲 partition 导致的Kafka 某些分区的消息 key 分布不均匀标签数据大部分落在少数几个分区上这几个分区的消费速度被 Sink 限制其他分区没有消息看起来 lag 不涨但实际消费进度被卡住。这个问题在 Flink 侧没有太好的自动解决手段只能从源头优化 Kafka 的 key 设计。我们最终把 Kafka 消息的 key 从纯粹的 user_id 改成了user_id % 分区数让标签事件更均匀地散布到所有分区。改造后同样作业的消费延迟下降了 60% 以上。5.4 OOM 问题RocksDB 状态后端是真的救星标签计算作业里如果有 Flink SQL 的 window 聚合会用到状态存储。双写作业上线后的第一个月TaskManager 偶发 OOM排查到最后是默认的 Heap 状态后端把窗口聚合的中间结果全部压在堆内存里。我当时的修改是启用 RocksDB 状态后端同时做了三个关键配置state.backend.rocksdb.memory.managedtrue让 Flink 自动管理 RocksDB 内存避免和 JVM 堆内存互相争抢。state.backend.rocksdb.block.cache.size128MBblock cache 不设置的话默认很小读性能会很差。state.backend.rocksdb.writebuffer.count4默认值是 2在标签高频更新场景下频繁 flush 会导致 CPU 飙升调大一点能让写入更平滑。RocksDB 切完之后TM 堆内存稳定下来了GC 停顿也明显变少。代价是本地磁盘多占用了一些状态存储空间但对生产环境来说完全可接受。如果你的集群条件允许建议所有实时标签类作业统一启用 RocksDB 状态后端不要等出了 OOM 再改。5.5 双写数据校验用对账任务兜底最后分享一个我认为双写作业务必加上的东西对账任务。很多人觉得 Flink 作业已经做了幂等就不需要再关心数据一致性了但分布式系统里最怕的是“逻辑正确但物理不一致”——比如 StarRocks 导入成功了但 Flink 又重试了一次导致某个版本被覆盖成旧数据。我的做法是写一个独立的 Flink 批作业每天凌晨扫描两边的(user_id, tag_date)明细数据对比行数和tag_json的 hash。不一致的数据输出到一个 Kafka topic由修复作业按小时回放。这个对账任务本身不复杂但它能兜住日常开发中 90% 以上的“边缘 case”值得花两天时间做出来。排障笔记之外再分享一点双写作业的性能调优没有一劳永逸的方案业务数据分布变了、标签字段变了、存储集群扩容了都可能导致原来最优的参数变得不再合适。我养成了一个习惯每次对作业做了配置变更后都要对比变更前后 7 天的 checkpoint 时长、反压比例和 Kafka lag 指标而不是只看当天的数据。数据积累久了哪些参数需要跟着业务周期调整心里就有底了。