You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 23:52:40