- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
Transform 文档通常是"单个 Transform 插件"的参数说明,但实际落地时,用户更需要能直接复制的"场景级 recipe",把下面三件事串起来:Source 如何把原始数据解析成一行(Row)、Transform 如何抽取/清洗/转换字段、Sink 对 Schema/类型的要求(尤其是向量库)。本文基于 docs/zh/transforms/recipes.md 展开,结合 SeaTunnel 仓库中的 Transform V2 源码实现,给出两个可直接改造的完整作业示例(深层嵌套 JSON 入 MySQL、MySQL 文本向量化写入 Milvus),并说明 LLM / 文件内容的处理边界。读完本文,你将掌握 JsonPath、Embedding、DynamicCompile 三个 Transform 的核心参数与底层行为,能独立拼装属于自己的场景化作业配置。
为什么需要"场景级 Recipe"
SeaTunnel 的作业是env、source、transform、sink四段式 HOCON 配置。单个 Transform 的文档只会回答"这个插件有哪些参数",却不会告诉你:Kafka 的 message value 怎么变成一行可抽取的列?MySQL 里的 JSON 字符串如何变成 Milvus 需要的FLOAT_VECTOR?这类问题必须放在完整链路中才能讲清楚,这也是"Recipe"存在的意义——把 Source、Transform、Sink 三者对 Schema/类型的约束一次性对齐。
仓库中 Transform V2 的插件实现位于 seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform,下文涉及的 JsonPath、Embedding、DynamicCompile 都可以在该目录下找到对应源码。
Recipe 1:解析深层嵌套 JSON(Kafka/File)并写入 MySQL
适用场景
- 上游数据是一个嵌套较深的 JSON(Kafka message value 或 JSON lines 文件)。
- 你只需要把其中少量字段抽出来,映射为"扁平列",写入关系型数据库表。
关键点:把原始 JSON 保留成单列
JsonPath transform 是从某一个字段(STRING/BYTES)里读取 JSON 并做 JsonPath 抽取的,因此链路的第一步是让"整段 JSON"成为一个独立列:
- Kafka:通常意味着你需要把 message value 作为一个字段保留下来(例如
STRING/BYTES)。 - File:可以用
LocalFile的text格式,让每行内容落到默认的content列上。
Kafka 示例:把 message value 保留成单列
如果你希望对 Kafka 的 message value 使用JsonPath,请确保 Kafka Source 能把"整段 message value"输出成单列。一个简单做法是使用format = text且不配置schema,这样 Kafka 会输出一个名为content的单列字符串字段:
source { Kafka { plugin_output = "kafka_raw" topic = "events" bootstrap.servers = "kafka:9092" format = text } } transform { JsonPath { plugin_input = "kafka_raw" plugin_output = "kafka_flat" columns = [ { src_field = "content" path = "$.event_id" dest_field = "event_id" dest_type = "string" } ] } }注意:如果你使用format = NATIVE,message value 字段名为value(bytes),此时src_field应配置为value。
源码视角:JsonPath 支持哪些输入类型
从 JsonPathTransform.java 的实现看,src_field指向的字段除了STRING和BYTES(字节数组直接new String(...)还原成 JSON 文本),还支持ARRAY、MAP、ROW类型——后三者会先经过JsonUtils序列化成 JSON 字符串再走 JsonPath 抽取。也就是说,即便上游解析出了结构化的 Row/Map,JsonPath 依然可以基于序列化后的文本做路径抽取。另外,编译后的JsonPath对象存放在ConcurrentHashMap缓存中(JSON_PATH_CACHE.computeIfAbsent(path, JsonPath::compile)),同一个路径只会编译一次,批量数据下避免重复编译开销。
完整示例(LocalFile JSON lines -> JsonPath -> JDBC sink)
假设输入文件每行一个 JSON 对象,例如:
{"event_id":"e1","payload":{"user":{"id":1001,"name":"alice"},"order":{"amount":12.34,"currency":"USD"}}}env { job.mode = "BATCH" parallelism = 1 } source { LocalFile { plugin_output = "events_raw" path = "/data/events.jsonl" file_format_type = "text" } } transform { JsonPath { plugin_input = "events_raw" plugin_output = "events_flat" row_error_handle_way = "SKIP" columns = [ { src_field = "content" path = "$.event_id" dest_field = "event_id" dest_type = "string" }, { src_field = "content" path = "$.payload.user.id" dest_field = "user_id" dest_type = "bigint" }, { src_field = "content" path = "$.payload.user.name" dest_field = "user_name" dest_type = "string" }, { src_field = "content" path = "$.payload.order.amount" dest_field = "order_amount" dest_type = "double" }, { src_field = "content" path = "$.payload.order.currency" dest_field = "currency" dest_type = "string" } ] } } sink { Jdbc { plugin_input = "events_flat" url = "jdbc:mysql://mysql:3306/demo" driver = "com.mysql.cj.jdbc.Driver" username = "root" password = "root" database = "demo" table = "events_flat" } }拿到这个模板后,你只需修改path的 JSONPath 表达式、dest_field的目标列名与dest_type的目标类型,再替换 JDBC 的连接参数即可复用到其他 JSON 入表场景。
JsonPath 配置参数详解
结合 JsonPathTransformConfig.java 的定义,columns数组中每个元素支持以下键:
| 参数 | 说明 | 默认值 |
|---|---|---|
src_field | JSON 源字段名,必须是上游 schema 中真实存在的列 | 必填 |
path | 用于从 JSON 中选取字段的 JSONPath,可以是单个字符串或字符串数组 | 必填 |
dest_field | 输出字段名,可以是单个字符串或字符串数组 | 必填 |
dest_type | 输出字段类型,可以是单个字符串或字符串数组 | string |
几个值得注意的源码级细节:
path、dest_field、dest_type均支持数组写法(例如一次抽取多个字段),解析后的数组长度必须一致,否则会抛出 "Path, dest_field, and dest_type arrays must have the same length" 异常(见 JsonPathTransformConfig.java)。src_field必须存在于输入 schema:配置解析阶段会校验table.getTableSchema().contains(srcField),找不到会抛出cannotFindInputFieldError,因此不要在 JSON 文本里凭空指定源列。dest_type由SeaTunnelDataTypeConvertorUtil解析,最终通过JsonToRowConverters把抽取结果按目标类型做转换(例如字符串"1001"转为bigint)。
错误处理:row_error_handle_way 与 column_error_handle_way
上面的示例用了row_error_handle_way = "SKIP"。该参数定义在 TransformCommonOptions.java,取值范围与行为如下:
row_error_handle_way(行级):可选FAIL(默认)、SKIP、ROUTE_TO_TABLE。FAIL时数据格式错误会阻塞作业并抛出异常;SKIP时跳过该行数据;ROUTE_TO_TABLE可将脏数据路由到row_error_handle_way.error_table指定的目标表。column_error_handle_way(列级,可在每个 column 内配置):可选FAIL、SKIP、SKIP_ROW。SKIP时该列出错则置空跳过;SKIP_ROW时该列出错则跳过整行。
在 JsonPathTransform.java 中可以看到:当 JsonPath 抽取抛出JsonPathException(例如路径不存在)时,如果列级错误处理允许跳过,则返回null不阻塞;否则抛出ErrorDataTransformException。这一机制保证了脏数据不会拖垮整条链路,适合生产环境的容错诉求。仓库测试 JsonPathTransformTest.java 中对路径抽取与错误处理分支均有覆盖用例。
1 行输入能否变成多行输出?
Transform V2 的 SPI 设计上同时支持:
- map(1 行输入 -> 1 行输出)
- flatMap(1 行输入 -> N 行输出)
但目前内置插件层面并没有一个"explode(把数组字段拆成多行)"transform。如果你的 JSON 里有items: [...]这类数组字段,并希望把每个 item 展开成独立行,通常有以下选择:
- 在上游预处理(例如 Flink/Spark SQL,或由上游生产端直接输出"一行一条目标记录"的 JSON)。
- 写一个自定义 Transform V2(flatMap)插件来实现 1:N 拆分。
- 如果只是"拆成多列(仍然 1:1)",可以优先用内置的
JsonPath、Split、Sql、FieldMapper等。
Recipe 2:把 MySQL 数据转换为 Milvus FLOAT_VECTOR
Milvus 侧的类型要求
Milvus 的FLOAT_VECTOR对应 SeaTunnel 的FLOAT_VECTOR逻辑类型(不是array<float>)。在 SeaTunnel 运行时,FLOAT_VECTOR的值类型是ByteBuffer。这一点可以在 VectorType.java 中确认:VectorType.VECTOR_FLOAT_TYPE的类型参数就是ByteBuffer,SQL 类型为SqlType.FLOAT_VECTOR;同文件还定义了SPARSE_FLOAT_VECTOR、BINARY_VECTOR、FLOAT16_VECTOR、BFLOAT16_VECTOR等向量变体,均标注为@Experimental(实验特性,使用需谨慎)。
方案 A(推荐):Embedding transform 直接对文本向量化
如果你有原始文本(或其他 Embedding 支持的模态),希望在链路中直接生成向量,推荐使用Embeddingtransform。它会新增一个FLOAT_VECTOR列,并根据模型响应自动设置向量维度。
enable_auto_id = true只有在上游 schema 已经包含主键信息时才会生效。这个示例里,id应该是源表的真实主键,且 JDBC source 需要把该主键元信息保留下来。如果你的query或后续 transform 丢失了主键信息,请提前创建好 Milvus collection,或在上游 schema 中显式保留主键。
env { job.mode = "BATCH" parallelism = 1 } source { Jdbc { plugin_output = "mysql_docs" url = "jdbc:mysql://mysql:3306/demo" driver = "com.mysql.cj.jdbc.Driver" username = "root" password = "root" query = "select id, content from demo.documents" } } transform { Embedding { plugin_input = "mysql_docs" plugin_output = "mysql_docs_with_vec" model_provider = "OPENAI" model = "text-embedding-3-small" api_key = "sk-xxx" vectorization_fields { content_vector = content } } } sink { Milvus { plugin_input = "mysql_docs_with_vec" url = "http://milvus:19530" token = "<your-token>" collection = "documents" enable_auto_id = true } }Embedding 的核心参数与底层行为
Embedding插件的实现位于 EmbeddingTransform.java,配置定义分布在 EmbeddingTransformConfig.java 与基类 ModelTransformConfig.java 中:
model_provider:从 EmbeddingTransform.java 的open()逻辑看,支持OPENAI、DOUBAO、QIANFAN、ZHIPU、AMAZON(Amazon Bedrock)、CUSTOM(通过custom_config自定义请求/响应解析)等提供方;LOCAL会直接抛出IllegalArgumentException。model:具体模型名,例如 OpenAI 的text-embedding-3-small。api_key/secret_key:QIANFAN、AMAZON等提供方还需要secret_key,AMAZON额外需要aws_region。api_path/oauth_path:自定义服务地址(api_path带openai.api_path回退键)。vectorization_fields:输入字段到输出向量列的映射关系(必填)。配置解析阶段会对每个输出列做inputRowType.indexOf(srcFieldName)校验,源字段不存在会抛cannotFindInputFieldsError。支持两种写法:字符串content_vector = content(默认文本模态),或对象形式{field = "content", modality = "image/jpeg", format = "url"}等多模态写法。single_vectorized_input_number:单次请求向量化的输入条数,默认 1(例如千帆接口限制单次最多 16 条消息时可调整)。process_batch_size:每批处理的行数,默认 100(回退键inference_batch_size)。model_retry_max_attempts:单次远程模型请求的最大尝试次数,默认 1 表示不自动重试;配合model_retry_backoff_ms、model_retry_max_backoff_ms、model_request_timeout_ms控制重试与超时。dimension:默认 2048,仅在部分提供方(如ZHIPU、AMAZON)显式使用;对OPENAI等提供方,维度由模型响应自动确定——open()中dimension = model.dimension()取自模型返回结果,随后输出列以VectorType.VECTOR_FLOAT_TYPE建列(见 EmbeddingTransform.java)。这正是文档所说"根据模型响应自动设置向量维度"的源码依据。
关于 enable_auto_id 的提醒
enable_auto_id是 Milvus Sink 的参数,定义在 MilvusSinkOptions.java,默认值为false。它会启用 Milvus 侧的自动主键,但不会代替你提供主键列——enable_auto_id = true只有在上游 schema 已经包含主键信息时才会生效。因此示例中 JDBC Source 必须把真实主键id查出来并保留其元信息;如果链路中主键信息丢失,请提前创建好 Milvus collection,或在上游 schema 中显式保留主键。
方案 B:MySQL JSON 数组 -> FLOAT_VECTOR(需要自定义转换)
如果你的 MySQL 表里已经存了 embedding(例如 JSON 字符串:"[0.12, 0.98, ...]"),仍需要把它转换成 SeaTunnel 的FLOAT_VECTOR(ByteBuffer)。目前内置插件层面没有一个专门的"JSON 数组转 FLOAT_VECTOR"的 transform。你可以选择:
- 写自定义 Transform V2 插件来做类型转换;
- 或使用
DynamicCompile在作业里内联代码完成解析与ByteBuffer构造。
这里同样适用主键规则:enable_auto_id = true不会自动创建主键。上游 schema 仍然需要提供真实主键列,例如id。
下面给出一个使用DynamicCompile(Groovy)的完整示例:从 MySQL 读取 JSON 字符串列embedding_json,并转换成FLOAT_VECTOR列content_vector(示例维度为 4,请按实际 embedding 维度调整)。
该示例通过列下标读取embedding_json(inputRow.getField(1)),因此请保持query列顺序为id, embedding_json,或自行调整下标。
env { job.mode = "BATCH" parallelism = 1 } source { Jdbc { plugin_output = "mysql_embeddings" url = "jdbc:mysql://mysql:3306/demo" driver = "com.mysql.cj.jdbc.Driver" username = "root" password = "root" query = "select id, embedding_json from demo.embeddings" } } transform { DynamicCompile { plugin_input = "mysql_embeddings" plugin_output = "mysql_embeddings_vec" compile_language = "GROOVY" compile_pattern = "SOURCE_CODE" source_code = """ import org.apache.seatunnel.api.table.catalog.Column import org.apache.seatunnel.api.table.catalog.CatalogTable import org.apache.seatunnel.api.table.catalog.PhysicalColumn import org.apache.seatunnel.api.table.type.SeaTunnelRowAccessor import org.apache.seatunnel.api.table.type.VectorType import org.apache.seatunnel.common.utils.JsonUtils import org.apache.seatunnel.common.utils.VectorUtils import java.util.List class demo { Integer dim = 4 public Column[] getInlineOutputColumns(CatalogTable inputCatalogTable) { PhysicalColumn vectorCol = PhysicalColumn.of("content_vector", VectorType.VECTOR_FLOAT_TYPE, null, dim, true, null, "") return new Column[] { vectorCol } } public Object[] getInlineOutputFieldValues(SeaTunnelRowAccessor inputRow) { Object json = inputRow.getField(1) if (json == null) { return new Object[] { null } } List<Float> list = JsonUtils.toList(json.toString(), Float.class) if (list.size() != dim) { throw new IllegalArgumentException("embedding dimension mismatch, expected " + dim + " but got " + list.size()) } Float[] floats = list.toArray(new Float[0]) return new Object[] { VectorUtils.toByteBuffer(floats) } } } """ } } sink { Milvus { plugin_input = "mysql_embeddings_vec" url = "http://milvus:19530" token = "<your-token>" collection = "embeddings" enable_auto_id = true } }DynamicCompile 参数说明
DynamicCompile的配置定义在 DynamicCompileTransformConfig.java:
| 参数 | 说明 | 默认值 |
|---|---|---|
source_code | 要编译的源码,必填 | 无 |
compile_language | 编译语言(枚举),必填 | 无 |
compile_pattern | 编译模式(枚举),可选SOURCE_CODE等 | SOURCE_CODE |
absolute_path | 源码的绝对路径(配合compile_pattern使用) | 无 |
理解这段 Groovy 代码的关键在于它实现了 Transform V2 的两个内联钩子:
getInlineOutputColumns(CatalogTable):声明输出列。示例用PhysicalColumn.of("content_vector", VectorType.VECTOR_FLOAT_TYPE, null, dim, true, null, "")声明一个维度为dim、类型为FLOAT_VECTOR的新列。getInlineOutputFieldValues(SeaTunnelRowAccessor):逐行计算输出值。示例用inputRow.getField(1)按下标取embedding_json,经JsonUtils.toList解析为List<Float>,校验维度后通过VectorUtils.toByteBuffer(floats)构造ByteBuffer——这正是 Milvus Sink 期望的FLOAT_VECTOR运行时值类型。
仓库测试 JsonPathTransformTest.java 展示了同类"配置化构造 Transform"的测试范式,可作为自写插件时对照参考。如果后续数据规模变大,更稳妥的做法仍是把这类转换固化为自定义 Transform V2 插件,避免在作业配置中维护内联代码。
LLM / 文件内容处理边界(内置支持 vs 需要前置处理)
LLMtransform 是把一行里的字段作为输入发送给模型推理,它不会自动读取文件路径、下载 URL,也不会自动对大文件做切片/分块。- 如果要处理大文件或二进制内容,建议在进入
LLMtransform 之前完成解析与切片(例如用文件类 connector 读成文本行,或在上游预先把内容切成多行)。
从架构上理解这条边界:Transform V2 处理的是"行内字段"级别的数据,文件系统的读取职责属于 Source connector 层。需要处理大文件时,应当先用LocalFile、S3File等文件类 Source 把内容读成文本行(或分块行),再交给LLMtransform 逐行推理;二进制内容同理,应在上游完成解码与分片。另外,从 EmbeddingTransform.java 的二进制处理逻辑可以推断,Embedding 侧对二进制输入已有[data, relativePath, partIndex]三列约定与分片组装缓存,属于少数内置了二进制编排的模型类 Transform,普通LLMtransform 并不具备同等能力。
如何把 Recipe 改造成自己的作业
所有 Recipe 都可以直接放进 SeaTunnel 的config目录下作为作业配置文件提交运行。仓库提供了两个可参考的作业骨架:config/v2.batch.config.template(批处理模板)与 config/v2.streaming.conf.template(流处理模板)。改造时只需把握三条主线:
- Source 侧:确认原始 JSON 是否保留为单列(Kafka 用
format = text得到content列,NATIVE 格式则是value列;LocalFile 用file_format_type = "text")。 - Transform 侧:先用 JsonPath 做字段抽取与扁平化,需要向量时叠加 Embedding(有原始文本)或 DynamicCompile(已有 embedding 数组)生成
FLOAT_VECTOR列。 - Sink 侧:JDBC 关注表结构与目标列类型匹配;Milvus 关注
enable_auto_id与主键来源、向量维度是否与 collection 的 schema 一致。
生产环境建议再补充row_error_handle_way(脏数据容错)与 Milvus Sink 的data_save_mode、enable_dynamic_field、enable_upsert等参数(见 MilvusSinkOptions.java,enable_auto_id默认false、data_save_mode默认APPEND_DATA),按目标库的写入语义做显式配置,避免默认行为与预期不符。
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel Transform Recipes: Flatten Nested JSON into MySQL and Vectorize into Milvus FLOAT_VECTOR
SeaTunnel Transform Recipes: Flatten Nested JSON into MySQL and Vectorize into M
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Milvus Sink 连接器全解析:向量数据写入 Milvus 与 Zilliz Cloud 的实战指南
SeaTunnel Milvus Sink 连接器全解析:向量数据写入 Milvus 与 Zilliz Cloud 的实战指南 本指南基于当前仓库中的 Milv
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Milvus 源连接器实战:从 Milvus 与 Zilliz Cloud 读取向量数据的完整指南
SeaTunnel Milvus 源连接器实战:从 Milvus 与 Zilliz Cloud 读取向量数据的完整指南 本文基于 SeaTunnel 开源仓库中
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考