ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

SeaTunnel FakeSource 虚拟数据源连接器全解:配置参数、数据生成原理与多版本能力演进

SeaTunnel FakeSource 虚拟数据源连接器全解:配置参数、数据生成原理与多版本能力演进 SeaTunnel FakeSource 虚拟数据源连接器全解配置参数、数据生成原理与多版本能力演进【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelFakeSource 是 Apache SeaTunnel 内置的虚拟数据源连接器它不连接任何真实系统而是根据用户声明的 schema 结构随机生成指定数量的测试数据广泛应用于连接器功能验证、类型转换测试与本地管道联调。本文以 docs/zh/connectors/changelog/connector-fake.md 变更日志为骨架结合 FakeSource 官方文档 与 connector-fake 模块源码系统讲解 FakeSource 的全部配置参数、底层生成机制以及从 2.2.0-beta 到 2.3.12 的能力演进脉络帮助读者把 FakeSource 从随便生成几行数据用成精确可控的测试数据工厂。FakeSource 是什么虚拟数据源的定位与适用场景FakeSource 是一个虚拟数据源Source 连接器它根据用户定义的 schema 数据结构随机生成指定数量的行数据。由于不依赖任何外部系统它天然适合以下场景类型转换测试验证 SeaTunnel 类型系统对各种数据类型的支持map、array、row、decimal、bytes、date、timestamp、向量等连接器新功能验证在无外部依赖的前提下快速验证 Sink 连接器、Transform 插件的行为本地管道验证跑通 source → transform → sink 完整链路确认作业配置正确基准与压力测试通过row.num、split.num、parallelism控制数据规模评估下游吞吐。支持的引擎FakeSource 基于引擎无关的 SeaTunnel Connector API 开发可运行于以下引擎见 FakeSource.mdApache Spark批/流Apache Flink批/流SeaTunnel Zeta内置引擎特性支持矩阵根据 FakeSource.md 与 Connector V2 功能简介FakeSource 的特性支持情况如下特性支持情况说明批处理batch✅有界数据读完作业结束流处理stream✅作业模式为 STREAMING 时无界输出精确一次exactly-once❌虚拟数据源无需外部精确一次语义列投影column projection✅实现SupportColumnProjection接口并行度parallelism❌文档标注不支持但实现了SupportParallelism实际从 FakeSource.java 看类同时实现了SupportParallelism与SupportColumnProjection并行执行由分片机制承担支持用户自定义分片❌分片由枚举器自动计算不支持用户自定义分片规则值得说明的是FakeSource 的FakeSourceSplitEnumerator枚举器会一次性计算出全部 split并通过assignSplit(splitId % currentParallelism, ...)把它们确定性地分配给各个 reader 子任务。虚拟数据生成不需要跨 reader 的共享状态但 split → reader 的路由集中在枚举器中完成而不是由每个子任务自行枚举——这正是其并行能力与不支持用户自定义分片并存的原因。数据源选项详解完整参数表与默认值FakeSource 的配置项定义在 FakeSourceOptions.java 中并通过 FakeConfig.java 的buildWithConfig完成解析与校验。以下参数表继承自官方文档并补充了源码中的默认值依据表结构与数据量控制名称类型必填默认值描述tables_configslist否-定义多个 FakeSource 表每个项都可以包含单个 FakeSource 支持的完整配置如独立的 schema 和 rows。与顶层schema二选一配置schemaconfig条件必填-定义 Schema 信息未配置tables_configs时必填。支持fields与columns两种声明方式详见 Schema 特性rowsconfig否-自定义输出行列表每个并行度都会输出这组数据而不是随机生成row.numint否5每个并行度生成的数据总行数源码ROW_NUM默认 5split.numint否1枚举器为每个并行度生成的分片数量split.read-intervalint否1读取器在两个分片读取之间的间隔时间毫秒map.sizeint否5生成的map类型的大小array.sizeint否5生成的array类型的大小bytes.lengthint否5生成的bytes类型的长度string.lengthint否5生成的string类型的长度自增主键名称类型必填默认值描述auto.increment.enabledboolean否false是否启用自增 ID 生成2.3.12 新增auto.increment.startlong否1自增 ID 的起始值仅在auto.increment.enabledtrue时生效字符串与各数值类型的生成模式每个类型都支持两种生成模式FakeMode枚举RANGE随机范围 /TEMPLATE模板随机选取默认均为range类型模式选项最小值选项默认最大值选项默认模板选项stringstring.fake.mode--string.templatetinyinttinyint.fake.modetinyint.min0tinyint.max127tinyint.templatesmallintsmallint.fake.modesmallint.min0smallint.max32767smallint.templateintint.fake.modeint.min0int.max0x7fffffffint.templatebigintbigint.fake.modebigint.min0bigint.max0x7fffffffffffffffbigint.templatefloatfloat.fake.modefloat.min0float.max0x1.fffffeP127float.templatedoubledouble.fake.modedouble.min0double.max0x1.fffffffffffffP1023double.template向量相关2.3.8 新增名称类型必填默认值描述vector.dimensionint否4生成的向量维度二进制向量除外binary.vector.dimensionint否8二进制向量维度源码要求必须是 8 的倍数vector.float.minfloat否0向量中 float 数据的最小值vector.float.maxfloat否0x1.fffffeP127向量中 float 数据的最大值通用选项名称描述common-options数据源插件通用参数如plugin_input/plugin_output、parallelism等详见 Source Common Options参数校验逻辑在 FakeConfig.java 中buildWithConfig对 tinyint/smallint/int/bigint/float/double 以及 vector.float 的 min/max 均做了范围校验例如tinyint.min必须满足0 tinyint.min 127int.min必须落在0与Integer.MAX_VALUE之间越界会抛出FakeConnectorException错误码ILLEGAL_ARGUMENT。这保证了生成的随机值不会超出目标数据类型的合法区间。底层原理枚举器、读取器与数据生成器的三层协作FakeSource 遵循 SeaTunnel Source API 的枚举器Enumerator— 分片Split— 读取器Reader模型核心类位于 connector-fake/source 包1. 分片枚举FakeSourceSplitEnumeratorFakeSourceSplitEnumerator.java 是数据规模控制的核心discoverySplits()一次性算出全部 split对每张表按rowNum / splitNum向上取整得到每个 split 的行数再按numReaders当前并行度对每个 reader 进行分片编号每个 split 携带tableId、起始行索引和行数分配时通过split.getSplitId() % currentParallelism()计算归属的 reader 并加入pendingSplitssnapshotState(checkpointId)保存assignedSplits作为检查点状态FakeSourceState.java配合restoreEnumerator实现故障恢复2.3.0-beta 的修复#3112正是针对恢复时枚举器重复分配 split 的问题恢复后先从assignedSplits中剔除已分配分片。2. 数据读取FakeSourceReaderFakeSourceReader.java 负责逐分片产出数据维护一个并发安全的DequeFakeSourceSplit分片队列pollNext()在持有 checkpoint lock 的情况下取分片并调用对应的FakeDataGenerator生成数据通过split.read-interval取所有表配置的最小值控制相邻两个分片之间的读取间隔模拟限速常量MAX_ROWS_PER_POLL 4096单次pollNext最多发射 4096 行。该限制是为了避免大分片一次性长时间持有 checkpoint lock从而阻塞 checkpoint/savepoint barrier 注入导致 stop-with-savepoint 卡死。3. 数据生成FakeDataGeneratorFakeDataGenerator.java 是数据产出的最终执行者两种产出路径随机生成randomRow()遍历CatalogTable的物理列对每列调用randomColumnValue(column)按SqlType分发到 map/array/row/数值/时间等各类型生成逻辑最终包装成带tableId的SeaTunnelRow自定义行generateCustomRows()将用户rows配置序列化后经JsonDeserializationSchema反序列化并设置RowKindINSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE。4. 多表解析MultipleTableFakeSourceConfigMultipleTableFakeSourceConfig.java 决定单表还是多表模式若配置了tables_configsTABLE_CONFIGS走parseFromConfigs()每个子配置独立构建FakeConfig否则走parseFromConfig()构建单个FakeConfig当表数量 1 时校验各表tableId必须唯一否则抛出IllegalArgumentException。数据生成模式实战range 与 template随机范围模式range默认模式下FakeSource 根据各类型的 min/max 区间生成随机值。例如FakeSource { row.num 5 tinyint.min 1 tinyint.max 9 smallint.min 10 smallint.max 19 int.min 20 int.max 29 bigint.min 30 bigint.max 39 float.min 40.0 float.max 43.0 double.min 44.0 double.max 47.0 schema { fields { c_string string c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double } } }每个并行度输出 5 行各数值字段的取值严格落在配置的区间内。模板模式template当需要从指定候选值中随机选取而非区间取值时将对应类型的fake.mode设为template并配置模板列表FakeSource { row.num 5 string.fake.mode template string.template [tyrantlucifer, hailin, kris, fanjia, zongwen, gaojun] tinyint.fake.mode template tinyint.template [1, 2, 3, 4, 5, 6, 7, 8, 9] int.fake.mode template int.template [20, 21, 22, 23, 24, 25, 26, 27, 28, 29] bigint.fake.mode template bigint.template [30, 31, 32, 33, 34, 35, 36, 37, 38, 39] float.fake.mode template float.template [40.0, 41.0, 42.0, 43.0] double.fake.mode template double.template [44.0, 45.0, 46.0, 47.0] schema { fields { c_string string c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double } } }从源码看模板的解析入口在 FakeConfig.java模板列表在FakeDataRandomUtils中被随机选取。2.3.5 的修复#6438 fix random from template not include the latest value issue曾修正模板随机取值未覆盖最后一个元素的问题。类型声明与完整 schema 示例FakeSource 支持声明 SeaTunnel 全类型体系包括嵌套结构。一个覆盖主流类型的 schema 如下引用自 FakeSource.mdsource { FakeSource { row.num 16 schema { fields { c_map mapstring, string c_array arrayint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_decimal decimal(30, 8) c_null null c_bytes bytes c_date date c_timestamp timestamp } } plugin_output fake } }其中plugin_output将该表注册为可被下游plugin_input引用的临时表详见 Source Common Options。自定义 rows精确控制每一行数据与变更操作类型当需要精确控制输出内容例如构造特定的 CDC 变更序列时使用rows选项。它支持fields按 schema 字段顺序的取值列表与kind行类型两个字段source { FakeSource { schema { fields { c_string string c_int int c_bigint bigint } } rows [ { kind INSERT, fields [1, A, 100] }, { kind UPDATE_BEFORE, fields [1, A, 100] }, { kind UPDATE_AFTER, fields [1, A_1, 100] }, { kind DELETE, fields [1, A_1, 100] } ] } }从 FakeDataGenerator.java 的实现看kind会被映射为 SeaTunnel 的RowKindRowKind.valueOf(rowData.getKind())因此可以精确模拟 INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE 四种变更事件——这让 FakeSource 可以直接用于验证下游 Sink 对 CDC changelog 事件的处理。bytes 类型使用 base64 编码由于 HOCON 规范的限制用户无法直接创建字节序列对象FakeSource 使用字符串来为bytes类型赋值。在上面的示例中bytes字段被赋值为bWlJWmo——这是通过base64编码的miIZj。因此为bytes类型字段赋值时请务必使用 base64 编码的字符串。时间类型默认值CURRENT_TIMESTAMP / CURRENT_TIME / CURRENT_DATE对于时间类型可以通过rows或 schema 的columns声明默认值为CURRENT_TIMESTAMP、CURRENT_TIME、CURRENT_DATE来获取当前时间2.3.8 Time supports default value 引入的能力schema { fields { pk_id bigint name string score int time1 timestamp time2 time time3 date } } # 使用 rows rows [ { kind INSERT fields [1, A, 100, CURRENT_TIMESTAMP, CURRENT_TIME, CURRENT_DATE] } ]也可以使用columns方式schema { columns [ { name book_publication_time, type timestamp, defaultValue 2024-09-12 15:45:30, comment 书籍出版时间 }, { name book_publication_time2, type timestamp, defaultValue CURRENT_TIMESTAMP, comment 书籍出版时间2 }, { name book_publication_time3, type time, defaultValue 15:45:30, comment 书籍出版时间3 }, { name book_publication_time4, type time, defaultValue CURRENT_TIME, comment 书籍出版时间4 }, { name book_publication_time5, type date, defaultValue 2024-09-12, comment 书籍出版时间5 }, { name book_publication_time6, type date, defaultValue CURRENT_DATE, comment 书籍出版时间6 } ] }该能力由 FakeDataGenerator.java 的getNewValueForField实现当字段值为CURRENT_TIME/CURRENT_DATE/CURRENT_TIMESTAMP时分别替换为LocalTime.now()/LocalDate.now()/LocalDateTime.now()的字符串形式对于TIMESTAMP_TZ类型则替换为OffsetDateTime.now()对应 2.3.9 的 timestamp with timezone offset 支持。多表生成tables_configs 与 table-namesFakeSource 支持在单个作业中产出多张表这在用一个作业同时验证多条管道时非常实用。方式一tables_configsFakeSource { tables_configs [ { row.num 16 schema { table test.table1 fields { c_string string, c_tinyint tinyint, c_int int } } }, { row.num 17 schema { table test.table2 fields { c_string string, c_tinyint tinyint, c_int int } } } ] }从 MultipleTableFakeSourceConfig.java 可以看出tables_configs中每个子项都会独立构建一个FakeConfig拥有各自的row.num、schema 与 rows且当表数量大于 1 时会对所有表的 tableId 做唯一性校验。该能力由 2.3.4 的 FakeSource support generate different CatalogTable for MultipleTable#5766引入。方式二table-names2.3.4 同时引入了table-names选项#5604配合统一 schema 声明多张同名结构的表source { FakeSource { table-names [test.table1, test.table2, test.table3] parallelism 1 schema { fields { name string age int } } } }schema 中的表标识符与键约束2.3.4 起 schema 支持配置table标识符#5628并支持在 schema 中声明primaryKey、constraintKey与column级属性#5564。例如schema { fields { id int name string age int } primaryKey { name pk columnNames [id] } }这些元数据会进入CatalogTable使下游 Sink如基于主键做 upsert 的 Doris/StarRocks能拿到完整的表结构信息。Schema 的完整声明语法见 Schema 特性。向量数据生成面向 AI/RAG 场景的测试数据自 2.3.8 起FakeSource 支持生成向量数据#7401 Fake Source support produce vector data#7446 update vectorType可用于向量数据库如 Milvus、Qdrant连接器与 embedding 管道的测试source { FakeSource { row.num 10 vector.dimension 4 binary.vector.dimension 8 schema { table simple_example columns [ { name book_id, type bigint, nullable false, defaultValue 0, comment 主键 ID }, { name book_intro_1, type binary_vector, columnScale 8, comment 向量 }, { name book_intro_2, type float16_vector, columnScale 4, comment 向量 }, { name book_intro_3, type bfloat16_vector, columnScale 4, comment 向量 }, { name book_intro_4, type sparse_float_vector, columnScale 4, comment 向量 } ] } } }要点支持binary_vector、float16_vector、bfloat16_vector、sparse_float_vector等向量类型columnScale用于在列级别覆盖全局的维度设置binary.vector.dimension源码要求是 8 的倍数默认 8普通向量的维度由vector.dimension控制默认 4vector.float.min/vector.float.max控制向量元素随机值的范围默认 0 与Float.MAX_VALUE。2.3.12 又新增了 Support vector series sql function#9765在 Transform-V2 的 SQL 函数层面补齐了向量系列函数的支持使 FakeSource 产出的向量数据可以被 SQL Transform 直接加工。自增主键auto.increment 的使用自 2.3.12 起#9505FakeSource 支持自动递增 ID用于生成不重复的主键值source { FakeSource { plugin_output fake auto.increment.enabled true auto.increment.start 1000 row.num 50000 schema { fields { id int name string age int } primaryKey { name pk columnNames [id] } } } }实现上AutoIncrementIdGenerator.java 使用AtomicLong从auto.increment.start默认 1开始getAndIncrement()生成全局递增 ID。该特性配合 schema 中的primaryKey声明非常适合验证主键去重类 Sink如 MySQL、Kudu和需要唯一键的 upsert 语义测试。变更日志解读FakeSource 能力演进时间线connector-fake.md 变更日志 完整记录了 FakeSource 从雏形到成熟的关键节点。按能力领域归类如下条目中的 PR/Issue 编号均出自该变更日志数据生成能力版本变更能力含义2.2.0-betasupport user-defined-schema and random data for fake-table#2406支持用户自定义 schema 与随机数据FakeSource 基础能力诞生2.2.0-betasupports direct definition of data values (row)#2839引入rows直接定义数据行2.2.0-betaFake date calculation error#2573修复日期计算错误2.3.1Improve fake connector#3932连接器整体改进2.3.1Optimizing Data Generation Strategies#4061优化数据生成策略参考 issue #40042.3.5fix random from template not include the latest value#6438修复模板随机取值遗漏最后一项2.3.8Fake supports column configuration#7503schema 支持columns方式声明2.3.8Time supports default value#7639时间类型支持CURRENT_TIMESTAMP/CURRENT_TIME/CURRENT_DATE2.3.9Improve memory usage when split size is large#7821大分片场景的内存占用优化分片与并行版本变更能力含义2.3.0-betaSupport multi splits#2974支持多分片配合并行度提升吞吐2.3.0-betasupports setting split rows and reading interval#3098支持配置每分片行数与读取间隔2.3.0-betafix duplicate splits when restoring#3112修复恢复时重复分配分片2.3.1add parallelism and column projection interface#3829引入并行度与列投影接口多表能力版本变更能力含义2.3.4Addtable-namesfrom FakeSource/Assert#5604引入table-names多表产出2.3.4Support config tableIdentifier for schema#5628schema 支持表标识符2.3.4Support config column/primaryKey/constraintKey in schema#5564schema 支持列、主键、约束键2.3.4FakeSource support generate different CatalogTable for MultipleTable#5766多表各自生成独立 CatalogTable2.3.9Unified tables_configs and table_list#8100统一tables_configs与table_list配置向量与自增版本变更能力含义2.3.8Fake Source support produce vector data#7401支持向量数据生成2.3.8update vectorType#7446向量类型更新2.3.12Support auto-increment id#9505支持自增 ID2.3.12Support vector series sql function#9765Transform-V2 支持向量系列 SQL 函数框架与工程质量公共 API 演进版本变更影响2.2.0-betaReplace SeaTunnelContext with JobContext#2706移除单例模式作业上下文重构2.2.0-betaRename SeatunnelSchema to SeaTunnelSchema#2538命名规范化2.3.0Add Fake TableSourceFactory#3345引入工厂模式统一插件发现2.3.0Unified exception for fake source connector#3520统一异常处理2.3.0Fix option rule about all connectors#3592修复选项规则2.3.1Refactoring schema parse#4157schema 解析重构2.3.4Introduce new error define rule#5793新错误码定义规范2.3.4Add default implement for SeaTunnelSource::getProducedType#5670默认类型产出实现2.3.5Support event listener for job#6419作业事件监听2.3.8Add event notify for all connector#7501连接器事件通知2.3.9Rename result_table_name/source_table_name to plugin_input/plugin_output#8072配置命名迁移旧名已废弃见 Source Common Options 警告2.3.9Support timestamp with timezone offset#8367支持带时区偏移的时间戳2.3.10Improve fake source options#8950FakeSource 选项改进2.3.10restruct connector common options#8634公共选项重构2.3.11Add check script for source/sink state class serialVersionUID#9118增加状态类 serialVersionUID 缺失检查脚本实践建议与典型用法1. 端到端验证管道FakeSource 最常见的搭档是 Assert 连接器 与 Console Sink用 FakeSource 生成数据经 Transform 加工后由 Assert 校验字段类型与取值无需任何外部系统即可完成数据正确性闭环验证。2. 控制数据规模单并行度数据量 row.num行全作业数据量 ≈row.num × parallelism行用split.num增加分片数、用split.read-interval控制限速模拟分布式读取行为2.3.9 起大row.num 大分片场景的内存占用已得到优化#7821但超大测试集仍建议配合多并行度使用。3. 精确构造测试用例需要固定取值用rowstemplate模式需要唯一主键用auto.increment.enabled需要 CDC 变更序列用rows的kind字段组合 INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE需要时间语义用CURRENT_TIMESTAMP等默认值。4. 阅读源码的入口想要深入理解 FakeSource 实现建议按以下顺序阅读 connector-fake 模块config/FakeSourceOptions.java— 全部选项定义与默认值config/FakeConfig.java— 配置解析与范围校验config/MultipleTableFakeSourceConfig.java— 单表/多表模式分发source/FakeSourceSplitEnumerator.java— 分片计算与分配source/FakeDataGenerator.javautils/FakeDataRandomUtils.java— 随机数据生成核心source/FakeSourceReader.java— 分片消费与限速。FakeSource 作为零依赖的测试数据源其能力从 2.2.0-beta 的用户自定义 schema 随机数据一路演进到 2.3.12 的自增主键 向量数据 多表 自定义行 时间默认值已经成为 SeaTunnel 生态中连接器验证与管道联调不可替代的基础设施组件。掌握本文的参数表、生成机制与演进脉络即可把 FakeSource 从测试占位符升级为精确可控的数据生成工具。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进