☰
Flink DataGen实战:实时计算测试数据生成与压测避坑指南
2026/10/1 3:25:22 网站建设 项目流程

聊一个几乎每个做实时计算的人都会遇到的问题:测试环境没有数据源,但你又必须验证作业逻辑对不对、性能扛不扛得住、边界情况到底处理没处理。拉生产数据回来不现实,手写一个Java的Source又太重。我最近半年做实时任务,Flink DataGen SQL Connector是出现频率最高的工具之一。它不依赖外部消息队列、不用额外引入jar包,一条建表DDL就能持续吐数据,本地造数、压测、边界数据构造都能覆盖。这篇文章就把我常用的造数套路、参数坑、逼真数据模板都过一遍,适合刚接触Flink SQL的初学者,也适合已经会用DataGen但想把它玩得更透的人。

1. DataGen连接器到底是什么:先跑通一个最简单的例子

1.1 为什么是DataGen而不是其他造数方案

很多人第一次接触Flink SQL时,面对“没有数据源”的第一反应是写一个自定义SourceFunction,或者部署一套Kafka再用脚本往里面灌数据。这两种方式我都用过,但都有明显的副作用:手写SourceFunction意味着要维护一段Java代码,改字段结构就得重新编译打包;Kafka脚本灌数据则多了一个外部依赖,环境准备本身就要花不少时间。

DataGen最舒服的地方是它内置在Flink中,SQL Client打开就能用。它本质上是一个“按配置生成数据”的Source连接器,你不需要关心数据从哪里来,只需要在DDL里声明字段结构和生成规则,Flink就会根据配置自动产生数据流。对于本地联调、接口冒烟、回归测试来说,这是成本最低的方案。

1.2 一条DDL搞定连续数据流

以最常见的模拟订单表为例,先建一张带DataGen约束的源表:

CREATE TABLE datagen_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), ts TIMESTAMP(3) ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '10', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '1000000', 'fields.user_id.min' = '1', 'fields.user_id.max' = '1000', 'fields.amount.min' = '10.00', 'fields.amount.max' = '9999.99' );

然后在SQL Client里执行:

SET 'sql-client.execution.result-mode' = 'tableau'; SELECT * FROM datagen_orders;

你会看到数据源源不断地输出,每秒10行,id从1开始递增,user_id和amount在指定范围内随机。这就是DataGen最基础的用法。

这里有几个参数需要注意:rows-per-second控制每秒生成多少行,默认值是10000;fields.<字段名>.kind指定字段生成方式,有random和sequence两种,默认是random;sequence模式必须配置start和end,它会从start开始按步长1递增,走到end之后重新循环。

1.3 DataGen支持的字段类型与默认行为

DataGen官方文档说明它支持绝大多数Flink SQL类型,但不同类型的配置方式不一样,我整理了一张速查表:

字段类型常用配置默认行为
BOOLEAN无随机生成true/false
INT / BIGINT / SMALLINT / TINYINTmin/max在范围内均匀随机
DECIMALmin/max在范围内生成随机小数
FLOAT / DOUBLEmin/max在范围内生成随机浮点数
CHAR / VARCHAR / STRINGlength生成固定长度的随机字符串
ARRAY / MAPlength生成指定大小的随机集合
TIMESTAMP / TIMESTAMP_LTZmin/max在时间范围内随机生成
ROW无递归生成子字段随机值

一个容易忽略的点:字符串类型的length控制的是字符长度,如果你不配置,默认生成的字符串长度会比较短,做字符串截断、长度校验这类逻辑时容易测不出问题。数值类型如果不配置min和max,可能生成0到类型最大值之间的值,对于某些以0作为特殊值的业务字段,需要主动把min配成非0。

2. 压测时要搞清楚的三个核心点:速率、并行度和背压

2.1 rows-per-second的“每并行任务”陷阱

我第一次用DataGen做压测时犯过一个错误:以为rows-per-second = 10000就是整个作业每秒产生10000行。结果在Flink UI上看到source的numRecordsOutPerSecond到了好几万,一度怀疑自己配置错了。

后来翻官方文档才确认,rows-per-second是每个并行子任务每秒生成的行数。也就是说,如果source的并行度是4,实际总吞吐大约是4万行/秒。这个语义对整个压测结果影响巨大:你想压的是单并发处理能力还是总吞吐,对应的配置完全不同。

验证方法很简单:在Flink UI里找到source operator,看它的numRecordsOutPerSecond指标。如果作业并行度是N,这个值大约是rows-per-second的N倍。确认这一点之后再去调参数,才不会对着一个数值瞎猜。

2.2 多并行度下如何估算和设置总吞吐

明确了“每并行任务”的语义之后,设置总吞吐就变成一个简单的乘法问题。假设我希望整个作业达到每秒20万行的生成速率,又希望source有4个并行子任务,那么rows-per-second应该设置为50000。

并行度设置方式以SQL Client为例:

SET 'parallelism.default' = '4';

这个配置会同时影响source、算子、sink的并行度。如果你只想单独调整source的并行度,可以给表设置scan.parallelism?在DataGen的WITH参数中没有直接的parallelism选项,但可以通过作业级并行度来控制。对于压测场景,我一般先调parallelism.default,再根据UI上的瓶颈表现单独调整下游算子并行度。

需要注意的是,本地模式(LocalEnvironment)或者单机模式下的资源配置是有限的,不要把rows-per-second随便设置成上亿。数据生成本身也要消耗CPU,source算子如果成了瓶颈,你测的就不再是下游处理能力,而是DataGen自己的生成上限。

2.3 压测中怎么判断瓶颈在source、算子还是sink

实时作业的性能问题90%出在背压上。使用DataGen做压测时,我习惯把链路拆成三段来看:source -> 业务算子 -> sink。

在Flink Web UI的BackPressure页面,可以看到每个算子处于Idle、Busy还是BackPressured状态。如果source是Idle,说明它被下游反压了,你设置的生成速率根本没有完全发出去;如果source是Busy但sink没有出现背压,说明source自身生成数据的速度已经达到上限;如果中间业务算子长期Busy且反压比例很高,那瓶颈就在计算逻辑本身。

有一点很多人会忽略:如果sink用的是JDBC这类需要外部交互的连接器,外部系统一旦变慢,背压会直接传导到source,导致DataGen的吞吐实际达不到你预设的值。这时候你压的不是数据源,而是整个下游链路的综合能力。做纯source能力测试时,我建议sink用print或者blackhole,避免外部依赖干扰指标。print还能顺带看数据长什么样,blackhole则完全丢弃数据,是更干净的压测选择。

一个比较标准的压测链路:

CREATE TABLE datagen_source (...); CREATE TABLE blackhole_sink ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2) ) WITH ( 'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT id, user_id, amount FROM datagen_source;

然后打开UI观察source的numRecordsOutPerSecond和算子的背压状态,就能快速判断当前配置下作业能跑多快、瓶颈卡在哪里。

2.4 压测时要注意检查点与状态大小

用DataGen做长时间压测时,还有一个隐形坑:如果作业启用了Checkpoint,而下游sink的幂等性做得不好,或者状态没有合理配置TTL,状态会不断膨胀,最终拖垮整个作业。DataGen产生的数据本身无状态,但如果你在作业里做了聚合、窗口、维表关联,状态大小会持续增长。

我习惯在压测前把状态后端的相关参数明确调好,比如RocksDB状态后端的内存上限、state.ttl配置等,而不是用默认值跑长任务。否则跑了一小时数据量大增后,作业突然因为状态超限挂掉,你可能误以为是DataGen或下游的问题,实际上是状态管理没配好。

3. 边界数据不是碰运气:sequence字段与计算列的正确打开方式

3.1 sequence字段:造有规律数据的关键

DataGen的randommode适合造“分布随机”的数据,但很多场景我们需要的是“有规律且可控”的数据。比如主键必须递增、事件时间必须单调递增,这时候就得用sequence。

看一个实际例子:

CREATE TABLE datagen_seq ( seq_id BIGINT, user_id BIGINT ) WITH ( 'connector' = 'datagen', 'fields.seq_id.kind' = 'sequence', 'fields.seq_id.start' = '1', 'fields.seq_id.end' = '100000000', 'fields.user_id.min' = '1', 'fields.user_id.max' = '10000' );

seq_id会从1开始一路递增,到达end后从头循环。如果你需要一个永远不重复的字段,把end设得足够大,比如超过测试周期的总数据量即可。

“主键用sequence、业务字段用random”是我最常用的搭配。原因很简单:很多sink(比如JDBC、Iceberg、Hudi)对主键冲突非常敏感,DataGen默认的random模式生成主键时重复率极高,一旦出现主键冲突,写入就会报错。用sequence生成主键可以从根源上消除这类问题。

3.2 制造合法边界值:最大最小值、阈值附近

做边界测试时,我们要的不是“随机值”,而是“恰好等于边界或无限接近边界”的值。DataGen的min/max配置可以直接生成边界值区域的数据,但更多时候我会用计算列来精确控制。

比如要验证大金额字段对DECIMAL的精度处理,可以构造一个“恰好等于10000.00”的字段:

CASE WHEN id % 100 = 0 THEN CAST(10000.00 AS DECIMAL(10, 2)) ELSE CAST(0.01 + RAND() * 9999.99 AS DECIMAL(10, 2)) END AS amount

这一句的意思是每100条数据里就有一条触发最大边界值,其余99条在小额区间随机。这样既包含了边界,又不会让边界数据占满全部数据流。

类似的边界构造还有:

边界场景SQL表达式用途
BIGINT最大值CAST(9223372036854775807 AS BIGINT)测试溢出与序列化
金额0值CAST(0 AS DECIMAL(10, 2))测试零值分支逻辑
空字符串''测试非空校验
超长字符串REPEAT('a', 1000)测试字段长度上限

要注意,这些表达式里用到的字段必须来自DataGen生成的物理列,计算列不能直接引用另一个计算列。我用下来最顺手的做法是:在DDL里先定义几个DataGen物理列作为“随机种子”,再用计算列基于它们加工。

3.3 在数据流中注入NULL与异常值

Flink DataGen自身没有直接配置“多少比例生成NULL”的参数,但用计算列可以很轻松地实现。比如:

CREATE TABLE datagen_nullable ( id BIGINT, raw_amount DOUBLE, amount AS IF(RAND() < 0.05, NULL, raw_amount) ) WITH ( 'connector' = 'datagen', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '1000000', 'fields.raw_amount.min' = '1.0', 'fields.raw_amount.max' = '1000.0' );

这样约有5%的数据amount是NULL,下游在做COALESCE、CASE WHEN amount IS NULL之类的逻辑时就能得到有效验证。

有些业务还会遇到“字符串字段中混入非法值”的情况,比如数字字段偶尔出现负数、日期字段出现格式错误。这些也可以用类似思路注入。比如让status字段在90%情况下是'SUCCESS',10%情况下是异常枚举:

CASE WHEN RAND() < 0.9 THEN 'SUCCESS' ELSE 'UNKNOWN' END AS status

边界数据的目的不是让作业跑得更快,而是让那些平时遇不到的分支逻辑有机会被触发。我见过太多测试环境“一切正常”、一上生产就翻车的例子,根本原因就是测试数据太干净,下游对边界和异常完全没有抗性。

3.4 用sequence构造递增事件时间,配合watermark验证

实时计算里最经典的场景就是窗口计算,而窗口计算对事件时间的有序性非常敏感。DataGen默认生成的TIMESTAMP是随机时间,直接用它在测试窗口逻辑时会发现窗口永远触发不了,或者触发结果完全不符合预期。

解决办法是用sequence生成一个递增的BIGINT时间戳,再通过计算列转成TIMESTAMP_LTZ:

CREATE TABLE datagen_watermark ( seq BIGINT, event_time AS TO_TIMESTAMP_LTZ(seq, 3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'datagen', 'fields.seq.kind' = 'sequence', 'fields.seq.start' = '1577808000000', 'fields.seq.end' = '1600000000000' );

这里seq的单位是毫秒,起始值对应2020-01-01 00:00:00。sequence递增保证event_time严格单调,配合watermark,就能稳定测试滚动窗口、滑动窗口、迟到数据等逻辑。如果需要测试乱序数据,就把seq改成random模式,让时间戳在范围内随机跳动,从而模拟乱序到达。

这个技巧我几乎每个窗口类作业都会用到。没有它,你只能手写一个带递增时间戳的数据集,那又回到了Java Source的老路。

4. 让测试数据“像真的”:从随机白噪声到业务分布的改造

4.1 官方random kind生成的只是“白噪声”

直接用DataGen默认参数生成的数据,最大的问题是“假”。user_id在1到1000之间均匀随机,amount在10到9999.99之间均匀随机,这跟真实业务数据的分布规律完全不一样。真实订单的金额有小额高频、大额低频的长尾特征,真实用户活跃度有明显的时段差异,真实商品分类有热门冷门之分。

如果只是验证作业能跑通,白噪声没问题;但要验证聚合结果的正确性、优化SQL的执行计划、判断结果是否符合业务常识,白噪声数据会误导你。比如你在做“订单金额超过1000的用户占比”这类指标,均匀分布的数据和真实分布的数据算出来的结果可能相差很大,而你根本不知道哪个才对。

所以我后来养成了一个习惯:先用DataGen把基础随机字段铺好,再用SQL计算列把“业务形态”加工出来。数据生成的随机性来自DataGen,业务规则用SQL表达式控制,两者各司其职。

4.2 手机号、邮箱、金额等字段的SQL构造方法

先看一个完整的“模拟用户表”例子:

CREATE TABLE datagen_users ( id BIGINT, r DOUBLE, name AS CONCAT('user_', CAST(id AS STRING)), phone AS CONCAT('13', CAST(100000000 + CAST(r * 900000000 AS BIGINT) AS STRING)), email AS CONCAT('user_', CAST(id AS STRING), '@example.com'), amount AS CAST(ROUND(50 + r * 950, 2) AS DECIMAL(10, 2)), status AS CASE WHEN r < 0.9 THEN 'active' ELSE 'inactive' END ) WITH ( 'connector' = 'datagen', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '100000000', 'fields.r.min' = '0', 'fields.r.max' = '1' );

这里我把r作为一个DataGen物理随机字段,取值范围0到1,所有计算列都基于它派生。这样做的最大好处是:整个表内同一条记录的各字段之间是协调的——r < 0.9时status是active,amount也在小范围区间内,而不是每个字段各自独立随机导致数据之间毫无关联。

手机号构造的逻辑值得解释一下:100000000 + CAST(r * 900000000 AS BIGINT)的结果范围是100000000到999999999,正好是9位数字,前面拼上'13'就得到一个11位手机号。这种写法避免了字符串截断和补零问题,是我试过多种写法之后最可靠的一种。

4.3 模拟长尾分布:用权重区间逼近真实流量

真实世界的订单金额、接口耗时、用户消费频次几乎都是长尾分布。Flink SQL里没有内置的幂律分布函数,但我们可以用分段权重来近似,在实际测试中效果已经足够。

思路很简单:把0到1的随机数分成几个区间,每个区间对应一个数值范围。比如模拟订单金额时,70%的订单是小额、25%是中额、5%是大额:

amount AS CASE WHEN r < 0.70 THEN CAST(ROUND(1 + RAND() * 49, 2) AS DECIMAL(10, 2)) WHEN r < 0.95 THEN CAST(ROUND(50 + RAND() * 250, 2) AS DECIMAL(10, 2)) ELSE CAST(ROUND(300 + RAND() * 4700, 2) AS DECIMAL(10, 2)) END

第70百分位对应1到50元的订单,第95百分位对应50到300元,剩下5%对应300到5000元。这个分布虽然不是严格的长尾曲线,但已经能反映出“大部分订单金额不高,少数大单拉高整体GMV”的业务特征。

同样的手法可以处理很多场景:

  • 接口耗时:latency AS CAST((1 - r) * 500 AS DECIMAL(10, 2)),让大部分请求耗时短、少数请求耗时长
  • 设备类型:device AS CASE WHEN r < 0.5 THEN 'ios' WHEN r < 0.8 THEN 'android' ELSE 'web' END
  • 商品类目:用多个区间的CASE映射到不同类目ID

关于CASE里多个分支各自调用RAND()导致分布概率略有偏差的问题,我实际用下来影响不大。真正的关键是有一个全局的随机种子字段r来控制“走哪个分支”,至于分支内部的随机范围,再单独调用RAND()反而能让数据更分散。

4.4 当SQL表达式不够用时的进阶路线

计算列能覆盖70%的造数需求,但有些场景它确实搞不定:嵌套JSON结构、字段之间复杂的业务依赖、需要引用外部规则表的关联逻辑。这时候有两个方向可以选。

第一个方向是写UDTF,在SQL里注册一个自定义函数,传入DataGen生成的几个基础字段,函数内部做复杂加工后输出多行。优点是逻辑可以复用、可以单测;缺点是需要写Java代码并打包上传。

第二个方向是直接写一个Java的DataGenerator Source,完全放弃DataGen连接器。这个方向的灵活性最高,但也就失去了SQL造数的轻量性。我个人建议先用计算列,发现不行的再上UDTF,不要一开始就跳到Java Source。多数情况下,计算列配合CASE和内置函数已经能解决绝大部分“像真数据”的诉求。

5. 从造数到消费:本地测试闭环与常见踩坑记录

5.1 造完数往哪送:print、Kafka、JDBC三条链路

DataGen造出来的数据总得有个去处,不同场景选择不同sink。

最简单的验证方式是print连接器,它会把每条数据打印到TaskManager的日志中。启动方式:

CREATE TABLE print_sink ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2) ) WITH ( 'connector' = 'print' ); INSERT INTO print_sink SELECT id, user_id, amount FROM datagen_orders;

跑起来之后去TaskManager的stdout日志里就能看到输出。适合冒烟验证和查看数据长什么样。

如果需要跟真实链路对齐,可以接Kafka sink:

CREATE TABLE kafka_sink ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2) ) WITH ( 'connector' = 'kafka', 'topic' = 'test_orders', 'properties.bootstrap.servers' = 'localhost:9092', 'format' = 'json' );

Kafka sink需要引入flink-sql-connector-kafka的jar包,这点跟DataGen不一样,记得放到Flink的lib目录。

JDBC sink则适合把测试数据灌进MySQL这类数据库,供其他团队联调使用。JDBC链路有一个高频坑:源表主键用random生成时,插入MySQL很容易报主键冲突。解决方式前面已经提过,源表主键一律用sequence。

5.2 我踩过的DataGen高频问题

用DataGen这一年多,我整理了几个出现频率最高的问题:

问题现象根本原因解决办法
rows-per-second设置没用,速率远超预期该参数是“每并行任务”速率,总速率还要乘以并行度在UI上按并行度换算,或降低并行度观察
主键冲突导致sink写入失败random模式生成的主键重复率高主键字段改用sequence模式
窗口作业不触发默认timestamp随机无序,水位线一直不动用sequence+计算列构造递增事件时间
字段配置不生效,数据范围不对fields.<字段名>.min里的字段名和DDL字段名不一致检查字段名拼写,注意大小写和下划线
print sink看不到数据数据打印在TaskManager日志,不是客户端控制台看TM的stdout日志

第二行那个“字段配置不生效”特别隐蔽,因为Flink不会报错,只是安静地忽略你写错了的配置项,你看到的随机范围完全不是预期。发现数据“不够边界”时,第一反应应该是检查字段名是否与DDL完全一致。

5.3 让造数过程可复现:把DDL当代码来管理

本地测试最怕的事情是:昨天还能复现的数据,今天怎么就不一样了。DataGen是随机数据源,如果不做约束,每次跑出来的数据内容都会不同,边界数据的出现时机也不固定。

我的做法是把所有DataGen建表语句收敛到一个datagen.sql文件里,并提交到仓库管理。每个表的WITH参数上方都用SQL注释写明:这张表模拟什么业务、rows-per-second想表达的总吞吐是多少、sequence的start和end为什么这么设置。这样同事拿到文件就能复现,不需要我口头解释参数含义。

另外,对于需要固定数据量的测试,我不会让作业无限跑下去,而是用外层脚本启动作业并设置超时时间,或者手动cancel作业。DataGen自身没有“生成N条后自动停止”的参数,这一点需要结合实际场景在作业管理层面处理。

5.4 一个小提醒:DataGen不是生产组件

DataGen最大的优点是“无依赖、开箱即用”,但这也意味着它没有持久化、没有数据回溯、没有消息队列的ACK语义。它适合开发测试、压测验证、CI管道,不适合作为生产业务的数据源。我见过有人把DataGen直接跑在线上的一个演示任务里,短期内没问题,可一旦业务要接入真实数据,切换成本全落在了代码改造上。

我的建议是:DataGen用在一个相对隔离的测试环境里,所有依赖它的作业都通过统一的源表抽象来访问数据,等真实数据源接入时只需要替换DDL的WITH参数,业务逻辑基本不用动。

最后分享两个小技巧

如果让我总结DataGen用得顺手的核心,就一句话:把它当成一个“可以按需编程的数据发生器”,先用sequence解决顺序和主键,再用random和计算列解决分布和形态,最后用合适的sink验证完整链路。

我每次建表都会在SQL注释里写清楚这个表的用途和rows-per-second的语义,避免下次压测再被“每并行任务”这个数字误导。另一个小技巧是保存一份常用字段的造数模板,包含手机号、邮箱、枚举、时间戳、金额等字段的通用表达式,下次新开项目直接复制改参数,不用从零写一遍DDL。DataGen这类的内置连接器平时不起眼,但真正用熟了之后,你会发现它是本地实时开发效率提升最大的工具之一。

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

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

立即咨询