Flink是否具备Spark中用于测试的Rate Source功能?
Flink 中类似 Spark Rate Source 的测试数据源
Flink 提供了DataGen Connector,专门用于生成无界测试数据流,功能和Spark的Rate Source完全对标,支持控制数据生成速率、自定义字段结构等,满足各类流处理测试场景。
DataStream API 用法(Scala)
对应你给出的Spark Scala代码,Flink可以这样实现每秒生成10条数据的无界流:
import org.apache.flink.api.common.typeinfo.BasicTypeInfo import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment import org.apache.flink.streaming.api.functions.source.datagen.DataGenSource import org.apache.flink.streaming.api.watermark.WatermarkStrategy val env = StreamExecutionEnvironment.getExecutionEnvironment // 生成每秒10条的长整型序列流 val dataStream = env.fromSource( DataGenSource.builder[Long]() .typeInfo(BasicTypeInfo.LONG_TYPE_INFO) .rate(10) // 每秒生成10条数据 .build(), WatermarkStrategy.noWatermarks(), "Rate-like DataGen Source" ) dataStream.print() env.execute()
Table API/SQL 用法
如果习惯用SQL风格,也可以通过DDL定义测试数据源:
CREATE TABLE rate_like_source ( id BIGINT, event_time TIMESTAMP(3) ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '10', -- 控制每秒生成的数据量 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '9223372036854775807', 'fields.event_time.kind' = 'timestamp', 'fields.event_time.start' = '2023-01-01 00:00:00' ); -- 读取并使用数据 SELECT id, event_time FROM rate_like_source;
额外特性
DataGen还支持更多自定义配置:
- 生成随机字符串、枚举值等多种数据类型
- 控制字段的取值范围或分布
- 调整并行度来提升数据生成速率
内容的提问来源于stack exchange,提问作者Pavel Orekhov
相关产品推荐
相关产品推荐

