ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Flink SQL 上线前压测最短闭环:Print Sink 验证正确性 + BlackHole Sink 压性能

Flink SQL 上线前压测最短闭环:Print Sink 验证正确性 + BlackHole Sink 压性能 做 Flink SQL 任务上线前最怕的就是两眼一抹黑直接丢到生产。数据算得对不对性能顶不顶得住等问题暴露出来往往已经晚了。我这两年一直用的最短闭环方案就是Print Sink先把结果正确性验证清楚再切到BlackHole Sink把引擎性能上限榨干。这套流程不需要复杂的测试环境也不需要专门写一堆自定义 Sink一条 SQL 就能在本地跑完。这篇实战复盘会把整个压测闭环、常用算子模板以及我踩过的坑全部整理出来标好参数和步骤按着做就能复现。1. 为什么建议用最短闭环做 Flink SQL 压测1.1 压测到底想验证什么做压测之前先想清楚目标。很多人上手就调参、拉并行度结果忙了半天根本不知道自己在验证什么。Flink SQL 任务上线前无非两件事第一算出来的结果对不对第二有没有算得快、撑得住。这两个目标对应的验证手段完全不同。要验证正确性就得看到真实的数据流和输出结果最好能逐条核对要测性能上限就得把 Sink 端的影响拿掉让引擎把所有精力都花在计算本身。如果混在一起做拿到指标反而不知道是哪个环节拖了后腿。我之前见过一个团队压测时把结果写到 MySQLSink 端一慢整个作业吞吐掉了一半他们以为是 SQL 写得有问题排查了半天才发现是下游数据库写入瓶颈。这就是典型的“没想清楚压测目标”把验证正确性和压性能上限搅成了一锅粥。正确做法就是把这两件事拆开先用最直观的方式确认逻辑没问题再用最“无干扰”的方式测引擎极限。1.2 为什么选 Print 和 BlackHole 这对组合Flink 内置的两个 ConnectorPrint Sink 会把每一条结果打印到 stdout你打开 TaskManager 日志就能看到完整数据BlackHole Sink 则把写入的数据直接丢弃相当于“无底洞”。一个负责看一个负责扔恰好对应了正确性和性能两个目标。用 Print 的时候不需要接外部存储也没有写库失败带来的干扰可以快速核对字段、null 值、乱序后的结果用 BlackHole 的时候Sink 几乎不占额外资源测出来的吞吐就是纯计算逻辑加状态管理的极限。你可以把它理解成跑步测试Print 是背着摄像头跑每一步都能回放BlackHole 是卸掉所有装备跑看的就是你的真实耐力。这对组合最大的价值是“切换成本极低”验证正确性时用 Print压性能时把 INSERT 目标改成 BlackHole剩下的 SQL 一个字都不用动。这样的闭环才能快速迭代否则每次压测都在搭环境和造数据时间全耗在准备工作上了。1.3 最短闭环到底省掉了什么这里说的最短闭环是指Source - 计算逻辑 - Sink三段式链路中Source 使用 DataGen 模拟数据Sink 使用 Print/BlackHole不依赖 Kafka、MySQL、HDFS 这些外部系统。整个链路只需要一个 Flink 实例加 SQL Client就能完成从正确性验证到性能摸底的全部工作。省掉的东西非常可观不用搭测试 Kafka不用建测试库表不用准备脱敏数据更不用为了一个临时验证写一个自定义 Sink。最重要的是“可复现”同一个 SQL 丢进去就能得到同样的结果换人也能快速接手。这套闭环跑出来的性能数据虽然在绝对数值上不能直接等于生产但用来横向对比不同 SQL 写法、评估状态大小、定位算子热点完全够用。2. 最短闭环环境搭建与前置准备2.1 版本选择与依赖准备建议直接使用 Flink 1.17 或 1.18 的二进制包DataGen、Print、BlackHole 这几个连接器都已经内置不需要额外引 Jar。如果你还在用 Flink 1.13 之类的老版本大概率要自己补flink-table-planner相关依赖折腾成本高不如直接换新版本。安装步骤很简单下载解压后bin/start-cluster.sh拉起集群再bin/sql-client.sh进入客户端就能开始写 SQL 了。这种“单机起集群”的方式对压测来说完全足够因为我们关心的不是集群规模而是算子层面的计算逻辑和状态行为。如果你要模拟生产环境的资源配置记得把并行度、内存、State Backend 等参数按照生产规格提前配好否则本地压出来的数值没有参考意义。我在实践里习惯准备一个init.sql把所有公共 SET 参数放在里面进入 sql-client 后先执行一遍避免每次重复敲。2.2 用 DataGen 模拟高吞吐数据压测数据不该依赖真实业务库否则你无法控制流量还可能把上游数据库压垮。DataGen Connector 是 Flink 自带的模拟数据源可以在不连接外部系统的情况下持续生成指定速率的数据。它的rows-per-second参数直接控制每秒生成条数非常像把水龙头开到固定大小。配字段类型时可以指定sequence递增序或random随机值甚至可以给字段指定 min/max 范围在完全可控的数据源上构造出接近真实业务的字段分布。下面是我常用的订单源表模板。注意ts字段设置了 watermark这是后面做窗口聚合和 Interval Join 的前提如果源表没有事件时间和水位线正确性验证无从谈起。CREATE TABLE orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector datagen, rows-per-second 10000, fields.order_id.kind sequence, fields.order_id.start 1, fields.order_id.end 1000000000, fields.user_id.kind random, fields.user_id.min 1, fields.user_id.max 100000, fields.amount.kind random, fields.amount.min 1, fields.amount.max 10000, fields.ts.kind random, fields.ts.min 2024-01-01 00:00:00.000, fields.ts.max 2024-01-01 01:00:00.000 );这里有个细节ts用随机时间戳会制造大量乱序数据watermark 设成 5 秒延迟可以触发窗口和 Join 的计算。DataGen 是并行生成的每个子任务都会独立产生随机时间戳所以数据乱序程度比真实场景更夸张适合用来验证时间语义是否健壮。如果你希望 ts 尽量有序且时间向前推进可以用计算列生成当前时间戳但要注意所有并行子任务拿到的可能是同一个时间戳适合做语义验证不适合模拟真实乱序。2.3 一张 Print 表和一张 BlackHole 表Print 和 BlackHole 的建表语法非常简单。注意字段类型要和结果集对齐字段名无所谓但类型必须匹配否则会转换报错。我一般会把 Print 表和 BlackHole 表做成相同的结构这样在验证正确性和压测性能之间切换只需要改 INSERT 的目标表名。CREATE TABLE print_sink ( user_id BIGINT, user_name STRING, order_cnt BIGINT, total_amount DECIMAL(20, 2) ) WITH ( connector print ); CREATE TABLE blackhole_sink ( user_id BIGINT, user_name STRING, order_cnt BIGINT, total_amount DECIMAL(20, 2) ) WITH ( connector blackhole );这里提醒一句Print Sink 的结果会打印到对应子任务的 TaskManager stdout 日志里。如果你没看到输出不要急着说没数据先去看flink-dir/log目录下的taskmanager*.log或.out文件。并且 Print 默认会在每行前面加task-subtaskIndex前缀方便你分辨是哪条并行子任务输出的。这个序号分布也可以当作一个简单的并行度验证手段。3. 核心模板Join / Agg / TopN / UDF 的正确性与性能验证3.1 验证正确性INSERT INTO print_sink构造好源表和 sink 表之后把你要验证的业务 SQL 直接嵌到INSERT INTO print_sink后面。这样一旦任务跑起来TaskManager 日志里会持续打印符合条件的每一条结果。看起来简单但很多人会忽略一点Print 只能证明“在你的输入数据范围内”结果正确它不会像写库成功那样给出“最终一致”信号。所以验证正确性时要在源表里故意构造边界数据比如 null 字段、ts 突跳、user_id 匹配不上等然后看 Print 出来的结果是否符合预期。不要只看正常数据跑通就认为万事大吉。我之前验证一个窗口聚合任务正常数据跑了五分钟输出看起来都合理。后来故意往数据里塞了一条ts比当前水位晚很多的事件发现窗口结果完全没有包含这条数据才意识到是 watermark 设置得太激进导致部分事件永远不会触发窗口计算。这种事如果不提前在 Print 阶段暴露上线之后就是数据延迟和漏数事故。3.2 压测性能改一句话切到 BlackHole正确性验证通过后把 INSERT 的目标从print_sink改成blackhole_sink其余 SQL 一律不动。这一句话的切换就是整套最短闭环最爽的地方。因为 BlackHole 不保留数据你测出来的吞吐不再受下游消费能力限制可以直接看到 Flink 计算引擎在这个逻辑上的天花板。切过去之后建议先小并行度跑 10 分钟看稳定性再逐步把并行度拉上去。需要记住的是BlackHole 并不是零开销它仍然要经过序列化和网络传输只是非常轻。所以实际生产 Sink 的吞吐会比 BlackHole 低这个数值上限适合用来横向对比不同 SQL 写法的性能差异不是生产承诺值。我一般会在压测结束后记录三个数据BlackHole 峰值吞吐、稳定期吞吐、以及开启 Checkpoint 后的吞吐。这三个数能快速判断一个任务还有多少性能余量。3.3 四类高频业务模板直接抄因为下面的模板会同时用到orders和users两张源表这里把用户维表也补上模拟用户基础信息流。user_id用 sequence 生成保证和订单表能关联上user_name用 random 生成。CREATE TABLE users ( user_id BIGINT, user_name STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector datagen, rows-per-second 10000, fields.user_id.kind sequence, fields.user_id.start 1, fields.user_id.end 100000, fields.user_name.kind random, fields.user_name.length 10, fields.ts.kind random, fields.ts.min 2024-01-01 00:00:00.000, fields.ts.max 2024-01-01 01:00:00.000 );3.3.1 Regular Join 模板适用场景两路流按主键关联比如订单流和用户流实时补齐维度字段。注意 Regular Join 在 Flink 里是双流状态关联默认会保留两侧所有未匹配数据状态没有边界数据量一大就越跑越慢。所以压测时一定要重点看状态增长而不是只看当前吞吐。INSERT INTO blackhole_sink SELECT o.order_id, o.amount, u.user_name FROM orders AS o JOIN users AS u ON o.user_id u.user_id;如果跑了十几分钟吞吐持续往下掉优先怀疑是 Regular Join 状态膨胀。可以用 TTL 约束状态存活时间比如SET table.exec.state.ttl 1 h;。但设置 TTL 后超过存活时间的数据可能会丢掉匹配必须在业务容忍度内使用。Print 验证阶段可以把 TTL 调大甚至不设先保证正确性压测阶段再开 TTL观察状态收窄后对吞吐的影响。3.3.2 Interval Join 模板适用场景订单和用户两个事件流需要限制在时间窗口内关联比如用户注册后 10 分钟内的下单行为。Interval Join 必须基于事件时间并且两侧都要有 watermark。这种语法比 Regular Join 更加贴合“事件先后发生”的业务语义状态也能自动清理是实时关联里更推荐的做法。INSERT INTO blackhole_sink SELECT o.order_id, u.user_name, o.ts FROM orders AS o JOIN users AS u ON o.user_id u.user_id AND o.ts BETWEEN u.ts - INTERVAL 5 MINUTE AND u.ts INTERVAL 10 MINUTE;这里最容易踩的坑是ON 条件里只写了等值条件没写时间区间。Flink 不会报错但结果会退化成普通关联时间语义完全失效。上线前一定要在 Print 模式下多打一些样本看两条流的时间差是否落在你设定的区间内。还有一个经验是Interval Join 的区间不宜设得过大否则状态保留时间长内存和 RocksDB 压力成倍上升压测指标会很难看。3.3.3 窗口聚合模板窗口聚合是 Flink SQL 里最常见的场景实时统计、实时报表都离不开它。下面这个模板按 1 分钟滚动窗口统计每个用户的订单数和金额它会同时测试窗口计算、分组聚合和状态写入。INSERT INTO blackhole_sink SELECT user_id, COUNT(*) AS order_cnt, SUM(amount) AS total_amount FROM orders GROUP BY TUMBLE(ts, INTERVAL 1 MINUTE), user_id;注意这里用的是 Event Time窗口触发依赖 watermark。如果换成 Processing Time结果会掺入算子本地时间不适合用来做生产正确性验证。用 Event Time 的时候数据乱序会导致窗口延迟触发Print 验证阶段要多等一会再判断窗口结果有没有输出。压测窗口聚合时可以适当开启 MiniBatch 聚合比如设置table.exec.mini-batch.enabledtrue、table.exec.mini-batch.allow-latency5s和table.exec.mini-batch.size5000。MiniBatch 能显著降低聚合算子与下游的交互次数但会引入最高允许延迟业务上能接受再用。3.3.4 TopN 模板适用场景实时排行榜比如按订单量、GMV 取 TopN。TopN 在流式场景里是一个非常考验状态的算子因为它本质上是一个持续更新的排序索引。下面模板按每分钟窗口聚合再取订单量最大的前 10 个用户。INSERT INTO blackhole_sink SELECT user_id, order_cnt, rn FROM ( SELECT user_id, order_cnt, ROW_NUMBER() OVER (ORDER BY order_cnt DESC) AS rn FROM ( SELECT user_id, COUNT(*) AS order_cnt FROM orders GROUP BY user_id, TUMBLE(ts, INTERVAL 1 MINUTE) ) ) WHERE rn 10;TopN 很容易出现全局数据倾斜。如果没有分组条件所有数据都会进入同一个排序算子热点问题非常明显。业务上如果允许尽量用PARTITION BY user_id或按维度拆开让每个分组并行排序否则压测到高吞吐时背压几乎必然爆表。另外TopN 的输出是更新流Print 模式下会看到对不同user_id的反复更新日志这是正常现象不要当成数据重复。3.3.5 UDF 模板Flink SQL 的内置函数覆盖不了所有业务场景自定义 UDF 是绕不开的。在 SQL Client 里你需要把写好的 UDF 打成 Jar然后用ADD JAR /path/to/udf.jar;添加依赖再执行CREATE FUNCTION注册。下面是一个简单的 ScalarFunction 示例用来给用户名拼接后缀import org.apache.flink.table.functions.ScalarFunction; public class SuffixUDF extends ScalarFunction { public String eval(String input, String extra) { if (input null) { return null; } return input _ extra; } }SQL 侧注册和使用的模板如下CREATE FUNCTION suffix_udf AS com.example.SuffixUDF; INSERT INTO blackhole_sink SELECT user_id, suffix_udf(user_name, detail) AS user_detail FROM users;UDF 压测时最容易发现性能瓶颈在反序列化和字符串拼接尤其是把 UDF 放到热点 Key 上时。建议在 Print 阶段多检查 null 输入UDF 内部任意抛异常都会导致整个作业失败判空保护不能省。需要额外强调的是如果你的 UDF 里有外部 IO比如查 Redis、查 MySQL千万别在模拟源表场景下直接压测那测的是外部系统而不是 Flink。正确姿势是先 mock 掉外部调用再单独压 UDF 本身的 CPU 开销。4. 压测执行与结果解读4.1 有效压测的四个前提第一先预热。Flink 任务启动后State、连接池、异步线程都处于冷状态前几分钟的数据没有参考价值。我习惯在开压前先跑 10 分钟然后再把指标清零重新观察。第二数据量要可控。DataGen 的rows-per-second应该从低到高拉满而不是一上来就开 10 万每秒否则连 Source 都变成瓶颈你看到的只是“数据进不来”。第三并行度要和真实环境一致。如果你的生产环境是 32 并行度本地用 2 个并发跑出的数字没有参考意义。第四Checkpoint 必须开。很多任务压测时不开 Checkpoint数据都在内存里速度当然快但生产上 Checkpoint 一开吞吐立刻下降。所以压测 SQL 前面要加上这段SET execution.checkpointing.interval 30s; SET execution.checkpointing.mode EXACTLY_ONCE; SET parallelism.default 8;Checkpoint 间隔也要根据业务容忍度来调。间隔太短状态频繁快照性能损耗大间隔太长故障恢复时丢的数据多。我的经验是先从 30 秒起步压在稳定吞吐的 80% 以下再考虑调大。4.2 盯指标不是只看 Throughput吞吐只是第一层指标。同一套 SQL吞吐高但反压严重说明数据在算子间堆积波形不稳。Flink UI 的每个算子节点上会显示 Send 指标和 BackPressure 水位。如果某个算子的背压百分比长时间超过 80就要点进去看 CPU 和 GC。这里推荐用火焰图配合jstack抓取线程栈能快速定位是序列化、Hash 冲突还是状态序列化导致的 CPU 热点。不要一边看总吞吐一边自我安慰要细化到单个算子的数据流入流出速率差一个数量级就说明有瓶颈。比如我之前压一个带正则解析的 UDF 任务总吞吐只有 3 万条每秒怎么看都像是下游 Sink 的问题。后来抓火焰图才发现 UDF 内部的正则编译占了 60% 的 CPU优化成预编译后同样资源下吞吐直接翻倍。这类问题只看总吞吐根本发现不了。4.3 性能瓶颈定位与调优思路如果发现吞吐低于预期按下面的顺序排查Source 是否够快KeyBy 是否出现数据倾斜State 是否过大UDF 是否有资源争抢SQL 逻辑是否存在不必要的 shuffle。调优时先调并行度再开 MiniBatch最后看状态后端。RocksDB 的内存参数不要一上来就手动调通常设置成 managed memory 模式更省心让 Flink 统一管理。另一个容易被忽略的点是序列化格式如果源表字段是 STRING 类型Flink 需要反复做字符串转 Timestamp、转 Decimal这些隐式转换会吃掉大量 CPU。尽量在上游就把字段类型定义对能省很多无谓开销。还要注意并行度不是越大越好。并行度上去之后网络 shuffle 和状态备份的开销也同步增加。我试过一个简单过滤任务并行度从 8 调到 16吞吐没有翻倍反而下降了 10%因为每个并行子任务都要做 checkpoint 上传状态小任务反而被快照开销拖住了。压测时并行度要阶梯式尝试记录每个并行度下的稳定吞吐而不是盲目拉满。5. 常见问题与排查技巧实录5.1 一跑 Print 任务日志被刷爆了Print 适合验证不适合长时间跑。如果你开着 10 万 rows/s 的数据源往 Print 里灌几分钟就能把磁盘撑爆。正确做法是验证阶段把rows-per-second降到 100 到 1000保证能看到每一类结果需要大数据量验证时输出到一个过滤后的精简结果集或者只打印 TopN 和聚合值。日志文件建议单独配置 log4j把 TaskManager 的 stdout 重定向到独立文件方便检索也不会和系统日志混在一起导致文件乱掉。另外一个实用技巧是在 Print 模式下可以给 SELECT 加LIMIT吗流式查询里LIMIT只在批式语义下有效流式任务直接加会报错或行为不符合预期。想控制输出量应该从 DataGen 的 rows-per-second 和总数据范围入手从源头控制而不是在结果侧截断。5.2 Print 能看到数据但切到 BlackHole 吞吐反而很低这个问题我遇到过好几次。原因通常是 Sink 端的并行度和 Source 端不一致或者某个算子中出现了数据倾斜。Print 阶段因为数据量小倾斜被掩盖了切到大数据量后倾斜立刻显现。排查方式很简单在 SQL 里给倾斜字段加PARTITION BY或改成GROUP BY后观察各子任务指标重点看有没有某个子任务长期满负载其他子任务闲置。如果确认倾斜先看字段取值分布再考虑加盐、二次聚合等方案。比如普通订单关联用户大多数订单都集中在少数热门用户上如果不处理Window 聚合和 TopN 都会热点集中。这时候可以先把 user_id 做一层随机拆分比如user_id % 100作为临时分组维度聚合完成后再合并能把热点打散不少。5.3 BlackHole 压测结果和线上实际差距很大BlackHole 测的是纯计算上限线上有真实 Sink、外部存储、网络抖动吞吐必然下降。建议再准备一个真实 Sink 作为参考端比如写入 Kafka 或 MySQL测出带下游的端到端吞吐和 BlackHole 结果做对比就能算出 Sink 环节消耗了多少性能。但真实 Sink 要有节流能力别把下游存储打爆。这个“上不封顶”的场景BlackHole 是最安全的真实 Sink 需要你心里有数。还有一种常见情况是BlackHole 压测时吞吐很漂亮切到生产后却有大量背压问题往往出在 Source 端不是 DataGen 而是 Kafka。DataGen 没有网络拉取的延迟和分区倾斜两者的 Source 行为差异不小。所以最短闭环适合验证计算逻辑和算子调优但在接近上线前建议用真实 Source 再跑一轮混合压测。5.4 问题速查表现象可能原因处理建议Print 无输出watermark 未触发调小水位线延迟或检查数据时间范围吞吐一直上不去Source rows-per-second 太低提高 DataGen 速率并调大 Source 并行度BackPressure 高单算子热点或 State 过大看火焰图定位调 KeyBy 策略或 TTLCheckpoint 失败RocksDB 占用过高或磁盘慢开 managed memory换 SSDUDF 报错输入包含 nullPrint 阶段提前检查并做判空保护切到 BlackHole 后性能下降数据倾斜加盐、二次聚合或调整并行度6. 最后的实操心得整套最短闭环看起来不复杂但能坚持执行的人不多。我的习惯是每次上线前先 Print 跑 20 分钟校验业务口径再 BlackHole 跑 1 小时记录峰值吞吐和反压曲线作为该 SQL 的“体检报告”留档。压测脚本和模板固化到项目仓库里新人来了照着跑一遍就能快速上手。最后提醒一句Print 和 BlackHole 再方便也不能替代你对自己业务口径的理解该核对的数据样本一定不要省。希望这套方法能让你下一次上线前心里更有底。
返回列表