SeaTunnel FakeSource 虚拟数据源连接器全解:配置参数、数据生成原理与多版本能力演进
2026/9/18 20:10:45 网站建设 项目流程

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

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

FakeSource 是 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.numsplit.numparallelism控制数据规模,评估下游吞吐。

支持的引擎

FakeSource 基于引擎无关的 SeaTunnel Connector API 开发,可运行于以下引擎(见 FakeSource.md):

  • Apache Spark(批/流)
  • Apache Flink(批/流)
  • SeaTunnel Zeta(内置引擎)

特性支持矩阵

根据 FakeSource.md 与 Connector V2 功能简介,FakeSource 的特性支持情况如下:

特性支持情况说明
批处理(batch)有界数据,读完作业结束
流处理(stream)作业模式为 STREAMING 时无界输出
精确一次(exactly-once)虚拟数据源无需外部精确一次语义
列投影(column projection)实现SupportColumnProjection接口
并行度(parallelism)❌(文档标注不支持,但实现了SupportParallelism实际从 FakeSource.java 看,类同时实现了SupportParallelismSupportColumnProjection,并行执行由分片机制承担
支持用户自定义分片分片由枚举器自动计算,不支持用户自定义分片规则

值得说明的是,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时必填。支持fieldscolumns两种声明方式,详见 Schema 特性
rowsconfig-自定义输出行列表,每个并行度都会输出这组数据(而不是随机生成)
row.numint5每个并行度生成的数据总行数(源码ROW_NUM默认 5)
split.numint1枚举器为每个并行度生成的分片数量
split.read-intervalint1读取器在两个分片读取之间的间隔时间(毫秒)
map.sizeint5生成的map类型的大小
array.sizeint5生成的array类型的大小
bytes.lengthint5生成的bytes类型的长度
string.lengthint5生成的string类型的长度

自增主键

名称类型必填默认值描述
auto.increment.enabledbooleanfalse是否启用自增 ID 生成(2.3.12 新增)
auto.increment.startlong1自增 ID 的起始值,仅在auto.increment.enabled=true时生效

字符串与各数值类型的生成模式

每个类型都支持两种生成模式(FakeMode枚举:RANGE随机范围 /TEMPLATE模板随机选取),默认均为range

类型模式选项最小值选项(默认)最大值选项(默认)模板选项
stringstring.fake.mode--string.template
tinyinttinyint.fake.modetinyint.min(0)tinyint.max(127)tinyint.template
smallintsmallint.fake.modesmallint.min(0)smallint.max(32767)smallint.template
intint.fake.modeint.min(0)int.max(0x7fffffff)int.template
bigintbigint.fake.modebigint.min(0)bigint.max(0x7fffffffffffffff)bigint.template
floatfloat.fake.modefloat.min(0)float.max(0x1.fffffeP+127)float.template
doubledouble.fake.modedouble.min(0)double.max(0x1.fffffffffffffP+1023)double.template

向量相关(2.3.8 新增)

名称类型必填默认值描述
vector.dimensionint4生成的向量维度(二进制向量除外)
binary.vector.dimensionint8二进制向量维度,源码要求必须是 8 的倍数
vector.float.minfloat0向量中 float 数据的最小值
vector.float.maxfloat0x1.fffffeP+127向量中 float 数据的最大值

通用选项

名称描述
common-options数据源插件通用参数,如plugin_input/plugin_outputparallelism等,详见 Source Common Options

参数校验逻辑

在 FakeConfig.java 中,buildWithConfig对 tinyint/smallint/int/bigint/float/double 以及 vector.float 的 min/max 均做了范围校验:例如tinyint.min必须满足0 <= tinyint.min <= 127int.min必须落在0Integer.MAX_VALUE之间,越界会抛出FakeConnectorException(错误码ILLEGAL_ARGUMENT)。这保证了生成的随机值不会超出目标数据类型的合法区间。

底层原理:枚举器、读取器与数据生成器的三层协作

FakeSource 遵循 SeaTunnel Source API 的"枚举器(Enumerator)— 分片(Split)— 读取器(Reader)"模型,核心类位于 connector-fake/source 包:

1. 分片枚举:FakeSourceSplitEnumerator

FakeSourceSplitEnumerator.java 是数据规模控制的核心:

  • discoverySplits()一次性算出全部 split:对每张表,按rowNum / splitNum向上取整得到每个 split 的行数,再按numReaders(当前并行度)对每个 reader 进行分片编号;
  • 每个 split 携带tableId、起始行索引和行数;
  • 分配时通过split.getSplitId() % currentParallelism()计算归属的 reader 并加入pendingSplits
  • snapshotState(checkpointId)保存assignedSplits作为检查点状态(FakeSourceState.java),配合restoreEnumerator实现故障恢复;2.3.0-beta 的修复(#3112)正是针对恢复时枚举器重复分配 split 的问题,恢复后先从assignedSplits中剔除已分配分片。

2. 数据读取:FakeSourceReader

FakeSourceReader.java 负责逐分片产出数据:

  • 维护一个并发安全的Deque<FakeSourceSplit>分片队列;
  • pollNext()在持有 checkpoint lock 的情况下取分片并调用对应的FakeDataGenerator生成数据;
  • 通过split.read-interval(取所有表配置的最小值)控制相邻两个分片之间的读取间隔,模拟限速;
  • 常量MAX_ROWS_PER_POLL = 4096:单次pollNext最多发射 4096 行。该限制是为了避免大分片一次性长时间持有 checkpoint lock,从而阻塞 checkpoint/savepoint barrier 注入,导致 stop-with-savepoint 卡死。

3. 数据生成:FakeDataGenerator

FakeDataGenerator.java 是数据产出的最终执行者,两种产出路径:

  • 随机生成randomRow()遍历CatalogTable的物理列,对每列调用randomColumnValue(column),按SqlType分发到 map/array/row/数值/时间等各类型生成逻辑,最终包装成带tableIdSeaTunnelRow
  • 自定义行generateCustomRows()将用户rows配置序列化后经JsonDeserializationSchema反序列化,并设置RowKind(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE)。

4. 多表解析:MultipleTableFakeSourceConfig

MultipleTableFakeSourceConfig.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.md):

source { FakeSource { row.num = 16 schema = { fields { c_map = "map<string, string>" c_array = "array<int>" 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_TIMESTAMPCURRENT_TIMECURRENT_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-names

FakeSource 支持在单个作业中产出多张表,这在"用一个作业同时验证多条管道"时非常实用。

方式一:tables_configs

FakeSource { 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-names

2.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 中声明primaryKeyconstraintKeycolumn级属性(#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_vectorfloat16_vectorbfloat16_vectorsparse_float_vector等向量类型;
  • columnScale用于在列级别覆盖全局的维度设置;
  • binary.vector.dimension源码要求是 8 的倍数(默认 8),普通向量的维度由vector.dimension控制(默认 4);
  • vector.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 起(#9505),FakeSource 支持自动递增 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 使用AtomicLongauto.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 #4004)
2.3.5fix random from template not include the latest value(#6438)修复模板随机取值遗漏最后一项
2.3.8Fake supports column configuration(#7503)schema 支持columns方式声明
2.3.8Time supports default value(#7639)时间类型支持CURRENT_TIMESTAMP/CURRENT_TIME/CURRENT_DATE
2.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(#5628)schema 支持表标识符
2.3.4Support config column/primaryKey/constraintKey in schema(#5564)schema 支持列、主键、约束键
2.3.4FakeSource support generate different CatalogTable for MultipleTable(#5766)多表各自生成独立 CatalogTable
2.3.9Unified tables_configs and table_list(#8100)统一tables_configstable_list配置

向量与自增

版本变更能力含义
2.3.8Fake Source support produce vector data(#7401)支持向量数据生成
2.3.8update vectorType(#7446)向量类型更新
2.3.12Support auto-increment id(#9505)支持自增 ID
2.3.12Support vector series sql function(#9765)Transform-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(#4157)schema 解析重构
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(#8950)FakeSource 选项改进
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. 精确构造测试用例

  • 需要固定取值用rows+template模式;
  • 需要唯一主键用auto.increment.enabled
  • 需要 CDC 变更序列用rowskind字段组合 INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE;
  • 需要时间语义用CURRENT_TIMESTAMP等默认值。

4. 阅读源码的入口

想要深入理解 FakeSource 实现,建议按以下顺序阅读 connector-fake 模块:

  1. config/FakeSourceOptions.java— 全部选项定义与默认值;
  2. config/FakeConfig.java— 配置解析与范围校验;
  3. config/MultipleTableFakeSourceConfig.java— 单表/多表模式分发;
  4. source/FakeSourceSplitEnumerator.java— 分片计算与分配;
  5. source/FakeDataGenerator.java+utils/FakeDataRandomUtils.java— 随机数据生成核心;
  6. 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),仅供参考

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

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

立即咨询