Flink DataGen Connector 深入解析:使用 DataGeneratorSource 生成测试数据流
Flink DataGen Connector 深入解析使用 DataGeneratorSource 生成测试数据流【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink导读DataGen Connector 是 Flink 内置的数据生成 Source它允许在不依赖 Kafka 等外部系统的情况下为 Flink 管道快速生成输入数据非常适合本地开发、Demo 演示和单元测试场景。本文将从 DataGen 的使用方式、并行切分原理、限速策略、有界性语义到精确一次保障等方面结合本仓库源码DataGeneratorSource.java 等进行完整剖析读者学完后可以熟练使用DataGeneratorSource构建任意数据类型的模拟数据流并准确控制数据速率与生成数量。一、DataGen Connector 概述DataGen Connector 为 Flink 管道提供了一种Source实现用于生成输入数据。它的典型价值在于本地开发或做 Demo 时无需搭建/连接外部系统如 KafkaConnector内置于 Flink无需额外引入依赖即可直接使用。从构建配置可以验证这一点flink-connector-datagen/pom.xml 中只声明了flink-core一个依赖且作用域为provided因此使用该 Connector 不会引入额外的传递依赖开箱即用。二、核心用法DataGeneratorSource 与 GeneratorFunctionDataGeneratorSource是 DataGen 的核心类它并行地产生 N 个数据点。其底层机制是将0到count-1的长整数序列切分成与并行子任务subtask数量相同的若干子序列向用户提供的GeneratorFunction逐个供应类型为Long的 index 值GeneratorFunction负责把子序列的Long值映射为任意数据类型的生成事件。源码中可见DataGeneratorSource内部组合了一个NumberSequenceSource参见 DataGeneratorSource.java通过new NumberSequenceSource(0, to)构造序列范围其中to count 0 ? count - 1 : 0——当count为 0 时退化为不产生任何元素的空 Source供 Table 内部测试使用。2.1 最小示例生成 1000 条记录以下代码将产生[Number: 0, Number: 2, ... , Number: 999]这样一条序列GeneratorFunctionLong, String generatorFunction index - Number: index; long numberOfRecords 1000; DataGeneratorSourceString source new DataGeneratorSource(generatorFunction, numberOfRecords, Types.STRING); DataStreamSourceString stream env.fromSource(source, WatermarkStrategy.noWatermarks(), Generator Source);GeneratorFunctionLong, String是一个输入为Long、输出为String的函数式接口其完整定义见 GeneratorFunction.java包含三个方法default void open(SourceReaderContext readerContext)初始化方法在真实数据映射前仅调用一次default void close()销毁tear-down方法O map(T value)核心映射逻辑将输入的Longindex 转换为输出元素。2.2 元素顺序与并行度元素的顺序取决于并行度每个子序列内部按序产生。因此当并行度限制为 1 时会产生一条从Number: 0到Number: 999的完全有序序列当并行度大于 1 时序列被拆分到多个并行子任务上各子序列内部有序但整体为多路数据流。这正是 DataGeneratorSource.java 类注释中描述的设计The source splits the sequence into as many parallel sub-sequences as there are parallel source readersSource 将序列切分为与并行 source reader 数量相同的子序列。2.3 从已有集合生成数据除自定义映射函数外仓库还在functions包中提供了两个内置的生成函数实现可作参考或复用FromElementsGeneratorFunction.java按序返回集合中的元素序列。它在map(Long nextIndex)中通过while (numElementsEmitted nextIndex)逻辑处理故障恢复时的位置对齐确保按 index 精确输出对应元素。IndexLookupGeneratorFunction.java基于 index 从集合中查表返回元素内部用TypeSerializer将元素序列化后缓存open()时反序列化构建lookupMapmap(index)直接返回lookupMap.get(index)。这两个实现有一个共同的约束可从 IndexLookupGeneratorFunction.java 的checkIterable看出集合中不允许出现 null 元素且所有元素必须是声明类型或其子类。另外若调用map()的次数超过集合元素个数即DataGeneratorSource的count设置大于集合长度会抛出NoSuchElementException提示应将产生记录数设置为与集合元素数相等——这是使用内置生成函数时最容易踩的坑。三、限速Rate Limiting控制数据产生速率DataGeneratorSource内置了限速支持可以在不牺牲真实性的前提下模拟不同吞吐的外部系统。3.1 按每秒记录数限速以下代码将以整体 Source 速率跨所有 Source 子任务求和不超过每秒 100 条的速度产生Long值流GeneratorFunctionLong, Long generatorFunction index - index; double recordsPerSecond 100; DataGeneratorSourceString source new DataGeneratorSource( generatorFunction, Long.MAX_VALUE, RateLimiterStrategy.perSecond(recordsPerSecond), Types.STRING);3.2 RateLimiterStrategy 提供的三种策略限速策略统一由RateLimiterStrategy工厂接口定义源码见 RateLimiterStrategy.java它提供了三个静态工厂方法策略工厂方法底层限速器说明按秒限速RateLimiterStrategy.perSecond(double recordsPerSecond)GuavaRateLimiter每个子任务分得recordsPerSecond / parallelism的配额总体速率不超过设定值实际产生数受并行拆分取整影响按检查点限速RateLimiterStrategy.perCheckpoint(int recordsPerCheckpoint)GatedRateLimiter限制每个检查点产生的记录数要求recordsPerCheckpoint parallelism否则会抛出IllegalArgumentException不限速RateLimiterStrategy.noOp()NoOpRateLimiter不限制记录速率是两参构造函数的默认策略RateLimiterStrategy实现了Serializable接口且DataGeneratorSource在构造时会通过ClosureCleaner.clean(...)对策略做递归闭包清理参见 DataGeneratorSource.java保证策略可被安全地序列化分发到各并行子任务。注意策略接口标注为Experimental后续版本 API 可能演进。四、有界性Boundedness语义DataGeneratorSource永远是有界的BOUNDED。这一点在源码中有直接证据getBoundedness()方法固定返回Boundedness.BOUNDED参见 DataGeneratorSource.java。但实际使用中存在一个伪无界技巧将count设为Long.MAX_VALUE从实践角度看序列永远不会结束从而等效于一个无界 Source这种用法非常适用于模拟持续不断的数据流配合第三节的限速策略即可模拟特定吞吐的实时数据源。对于有限序列官方文档建议考虑在BATCH执行模式下运行应用详见 execution_mode.md 中关于何时使用批执行模式的说明。切换批模式有两种方式# 通过命令行参数指定 bin/flink run -Dexecution.runtime-modeBATCH jarFile// 或在代码中显式设置 env.setRuntimeMode(RuntimeExecutionMode.BATCH);五、使用注意事项与一致性保证5.1 精确一次与至少一次保证DataGeneratorSource可以用于实现**至少一次at-least-once和端到端精确一次end-to-end exactly-once**的处理保证前提条件是GeneratorFunction的输出相对于其输入必须是确定性的deterministic——即相同的Long输入总是产生相同的输出。这是因为 Source 的故障恢复依赖对 index 序列的重放只有映射函数确定重放相同的 index 才能得到一致的结果进而配合 Flink 的检查点机制实现精确一次语义。从 FromElementsGeneratorFunction.java 的map()实现可以看到它在故障恢复时会根据nextIndex跳过已消费的元素位置这正是确定性 index → 确定性输出机制在底层落实的体现。5.2 在 Source 端直接产生确定性 Watermark利用 index 驱动的确定性生成机制还可以在 Source 端基于生成的事件和自定义WatermarkStrategy直接产生确定性的 Watermark。这意味着测试时不必依赖外部事件时间系统watermark 的推进与数据生成同样可预测、可复现非常适合验证窗口计算、乱序处理等逻辑。5.3 测试辅助工具仓库还提供了面向测试的辅助工厂类 TestDataGenerators.java其中fromDataWithSnapshotsLatch(...)可以创建一个先发出给定数据、等待两次检查点后再重发同样数据的特殊 Source底层结合IndexLookupGeneratorFunction与DoubleEmittingSourceReaderWithCheckpointsInBetween用于验证状态恢复、重复输出等场景说明 DataGen 在设计之初就充分考虑了与检查点/恢复机制的协同。六、总结DataGeneratorSource是 Flink 生态中一个小而精的测试利器其核心要点可归纳为零依赖内置无需额外引入外部系统与依赖即可生成数据index 驱动 函数映射通过GeneratorFunctionLong, OUT将Long序号映射为任意类型的数据天然支持确定性输出与确定性 watermark并行切分序列按并行度切分为子序列控制并行度即可控制数据的整体有序性内置限速RateLimiterStrategy.perSecond/perCheckpoint/noOp三种策略满足不同吞吐模拟需求有界语义 伪无界技巧BOUNDED是固定语义配合Long.MAX_VALUE与限速即可模拟持续数据流有限序列则建议使用BATCH模式运行一致性保障只要GeneratorFunction对相同输入产生相同输出即可支撑 at-least-once 与端到端 exactly-once 语义。无论是快速验证算子逻辑、编写集成测试还是在没有外部消息队列的环境下演示实时计算流程DataGen 都是值得优先考虑的数据源方案。建议读者进一步阅读本文引用的 DataGeneratorSource.java 与 RateLimiterStrategy.java以掌握更底层的实现细节。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

Ray Tune 自定义日志与日志产物:logging_example 实战解读

Ray Tune 自定义日志与日志产物:logging_example 实战解读

人工智能分布式训练强化学习任务调度模型推理服务 【免费下载链接】ray Ray is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads. 项目地址: https://gitcode.com/gh_mirrors/ra/ray 点…

2026/9/20 14:41:03 阅读更多 →
工业网络IBC密钥管理:SM9算法落地与工程实践踩坑实录

工业网络IBC密钥管理:SM9算法落地与工程实践踩坑实录

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 14:40:02 阅读更多 →
Steam本地manifest解析:Python一键导出游戏清单方案

Steam本地manifest解析:Python一键导出游戏清单方案

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 14:40:02 阅读更多 →

最新新闻

抖音视频批量下载与无水印保存完整指南:douyin-downloader 使用教程

抖音视频批量下载与无水印保存完整指南:douyin-downloader 使用教程

抖音视频批量下载与无水印保存完整指南:douyin-downloader 使用教程 【免费下载链接】douyin-downloader A practical Douyin downloader for both single-item and profile batch downloads, with progress display, retries, SQLite deduplication, and browser f…

2026/9/20 20:33:00 阅读更多 →
baoyu-wechat-summary 群友画像系统(Profiles)完整指南:文件格式、更新规则与回溯流程

baoyu-wechat-summary 群友画像系统(Profiles)完整指南:文件格式、更新规则与回溯流程

AI 技能AI 插件 【免费下载链接】baoyu-skills 项目地址: https://gitcode.com/gh_mirrors/ba/baoyu-skills 点击查看 免费下载 微信群的精华简报若要“越写越懂这群人”,靠的正是 per-user 画像(Profiles)体系。本文基于 baoyu-…

2026/9/20 20:33:00 阅读更多 →
PT 助手 Plus 安装:从源码到 Edge 商店的三条路径与避坑要点

PT 助手 Plus 安装:从源码到 Edge 商店的三条路径与避坑要点

PT 助手 Plus 安装:从源码到 Edge 商店的三条路径与避坑要点 【免费下载链接】PT-Plugin-Plus PT 助手 Plus,为 Microsoft Edge、Google Chrome、Firefox 浏览器插件(Web Extensions),主要用于辅助下载 PT 站的种子。 …

2026/9/20 20:33:00 阅读更多 →
GetQzonehistory:免费备份QQ空间历史说说、图片与留言

GetQzonehistory:免费备份QQ空间历史说说、图片与留言

GetQzonehistory:免费备份QQ空间历史说说、图片与留言 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 毕业前想把十年的空间说说和照片留下来,逐页翻看截图根本做…

2026/9/20 20:33:00 阅读更多 →
macOS 加载 SQLite-Vec 向量搜索扩展失败?三步编译加载的完整指南

macOS 加载 SQLite-Vec 向量搜索扩展失败?三步编译加载的完整指南

macOS 加载 SQLite-Vec 向量搜索扩展失败?三步编译加载的完整指南 【免费下载链接】sqlite-vec A vector search SQLite extension that runs anywhere! 项目地址: https://gitcode.com/GitHub_Trending/sq/sqlite-vec SQLite-Vec 向量搜索扩展在 macOS 上最…

2026/9/20 20:33:00 阅读更多 →
MATLAB实现物理信息神经网络PINN的故障诊断实战指南

MATLAB实现物理信息神经网络PINN的故障诊断实战指南

简介:面向工业故障诊断与智能运维场景,基于物理约束神经网络(PINN)的MATLAB分类预测实现,适合掌握一定编程和机器学习基础、希望将物理机理与数据驱动结合的研究人员、工程师和技术人员,可应用于航空航天、…

2026/9/20 20:32:00 阅读更多 →

日新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/20 0:00:46 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/20 0:00:46 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/20 0:00:46 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/20 0:00:46 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/20 0:00:46 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/20 0:00:46 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/19 23:01:36 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/19 17:50:38 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/19 23:35:34 阅读更多 →