Flink实时标签双写StarRocks与PolarDB-X的实践与调优
2026/9/14 20:08:05 网站建设 项目流程

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 是 OLAP,PolarDB-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,那么并行度设为 4~6 个消费者基本能打满存储端又不至于过度拉数据。

消费语义方面,双写场景我强烈建议用commit_on_checkpoint=true,也就是 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 是 bigint,tag_date 是日期字符串,这样既保证标签按天可回溯,也避免主键无限膨胀。StarRocks 侧再把 tag_date 设为分区字段,每天一个分区;PolarDB-X 侧把 user_id 作为分区键,保证同一个用户的标签查询落在同一个分片上。

字段层面,标签列我建议设计成“统一扩展字段 + 明细拆分字段”的组合。扩展字段是一个 map 类型或者 JSON 字符串,承接临时新增的标签,不用改表结构;明细字段是固定的高频标签,比如is_high_activebuy_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 Sink:Stream 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 Sink:JDBC 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 连接器需要把连接参数里加上rewriteBatchedStatements=true,同时在 SQL 初始化语句里写清楚主键更新逻辑。如果不加rewriteBatchedStatements,MySQL 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 的BackPressureCheckpoint指标一起看。当时我总结了一个判断流程:

  1. 如果Checkpoint Duration持续超过Checkpoint Interval,说明 Sink flush 已经阻塞主线。
  2. 打开 Flink Web UI 的 BackPressure 页面,定位到具体是哪个 Sink 子任务处于 High 状态。
  3. 看对应 TaskManager 的线程 dump,确认是阻塞在 HTTP 的 Stream Load 请求上还是 JDBC 的 executeBatch 上。
  4. 如果是 JDBC Sink,先看数据库侧有没有锁等待和慢查询;没有的话再调低sink.buffer-flush.max-rows并增大 Sink 并行度。

这套排查流程帮我定位了多次并非 Flink 本身问题、而是 PolarDB-X 一个二级索引导致插入慢的系统性问题。

注意:反压不是只能靠调大并行度硬扛。先确认反压是“下游存储能力不够”还是“单批写入效率太低”,两者的调优方向完全相反。前者要增加并行度,后者要增大批次大小。

5. 常见问题与排查技巧实录

5.1 Flink JDBC 连接器异常:Failed to upsert data的根源

这个报错在双写作业刚上线时几乎天天出现。排查下来发现是两类原因:

第一类原因是 JDBC 连接器默认的rewriteBatchedStatements=false。PolarDB-X 作为分布式数据库,驱动收到没有 rewrite 标记的 batch 请求时,实际上是逐条执行,每条数据一次网络往返。当批次数据量到几百上千条时,执行时间会被网络延迟放大,导致网关超时。解决办法很简单:在 JDBC URL 加上rewriteBatchedStatements=true&useServerPrepStmts=false

第二类原因是 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.managed=true:让 Flink 自动管理 RocksDB 内存,避免和 JVM 堆内存互相争抢。
  • state.backend.rocksdb.block.cache.size=128MB:block cache 不设置的话默认很小,读性能会很差。
  • state.backend.rocksdb.writebuffer.count=4:默认值是 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 指标,而不是只看当天的数据。数据积累久了,哪些参数需要跟着业务周期调整,心里就有底了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询