SeaTunnel Transform Recipes 实战指南:嵌套 JSON 解析、MySQL 向量化与 Milvus 入库
2026/9/20 11:31:07 网站建设 项目流程
  • 数据集成
  • ETL
  • 大数据
  • 批处理
  • 流处理
  • 变更数据捕获

【免费下载链接】seatunnel

SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/GitHub_Trending/se/seatunnel
点击查看免费下载

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 的作业是envsourcetransformsink四段式 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:可以用LocalFiletext格式,让每行内容落到默认的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指向的字段除了STRINGBYTES(字节数组直接new String(...)还原成 JSON 文本),还支持ARRAYMAPROW类型——后三者会先经过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_fieldJSON 源字段名,必须是上游 schema 中真实存在的列必填
path用于从 JSON 中选取字段的 JSONPath,可以是单个字符串或字符串数组必填
dest_field输出字段名,可以是单个字符串或字符串数组必填
dest_type输出字段类型,可以是单个字符串或字符串数组string

几个值得注意的源码级细节:

  • pathdest_fielddest_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_typeSeaTunnelDataTypeConvertorUtil解析,最终通过JsonToRowConverters把抽取结果按目标类型做转换(例如字符串"1001"转为bigint)。

错误处理:row_error_handle_way 与 column_error_handle_way

上面的示例用了row_error_handle_way = "SKIP"。该参数定义在 TransformCommonOptions.java,取值范围与行为如下:

  • row_error_handle_way(行级):可选FAIL(默认)、SKIPROUTE_TO_TABLEFAIL时数据格式错误会阻塞作业并抛出异常;SKIP时跳过该行数据;ROUTE_TO_TABLE可将脏数据路由到row_error_handle_way.error_table指定的目标表。
  • column_error_handle_way(列级,可在每个 column 内配置):可选FAILSKIPSKIP_ROWSKIP时该列出错则置空跳过;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)",可以优先用内置的JsonPathSplitSqlFieldMapper等。

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_VECTORBINARY_VECTORFLOAT16_VECTORBFLOAT16_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()逻辑看,支持OPENAIDOUBAOQIANFANZHIPUAMAZON(Amazon Bedrock)、CUSTOM(通过custom_config自定义请求/响应解析)等提供方;LOCAL会直接抛出IllegalArgumentException
  • model:具体模型名,例如 OpenAI 的text-embedding-3-small
  • api_key/secret_keyQIANFANAMAZON等提供方还需要secret_keyAMAZON额外需要aws_region
  • api_path/oauth_path:自定义服务地址(api_pathopenai.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_msmodel_retry_max_backoff_msmodel_request_timeout_ms控制重试与超时。
  • dimension:默认 2048,仅在部分提供方(如ZHIPUAMAZON)显式使用;对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_VECTORByteBuffer)。目前内置插件层面没有一个专门的"JSON 数组转 FLOAT_VECTOR"的 transform。你可以选择:

  • 写自定义 Transform V2 插件来做类型转换;
  • 或使用DynamicCompile在作业里内联代码完成解析与ByteBuffer构造。

这里同样适用主键规则:enable_auto_id = true不会自动创建主键。上游 schema 仍然需要提供真实主键列,例如id

下面给出一个使用DynamicCompile(Groovy)的完整示例:从 MySQL 读取 JSON 字符串列embedding_json,并转换成FLOAT_VECTORcontent_vector(示例维度为 4,请按实际 embedding 维度调整)。

该示例通过列下标读取embedding_jsoninputRow.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_CODESOURCE_CODE
absolute_path源码的绝对路径(配合compile_pattern使用)

理解这段 Groovy 代码的关键在于它实现了 Transform V2 的两个内联钩子:

  1. getInlineOutputColumns(CatalogTable):声明输出列。示例用PhysicalColumn.of("content_vector", VectorType.VECTOR_FLOAT_TYPE, null, dim, true, null, "")声明一个维度为dim、类型为FLOAT_VECTOR的新列。
  2. 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 层。需要处理大文件时,应当先用LocalFileS3File等文件类 Source 把内容读成文本行(或分块行),再交给LLMtransform 逐行推理;二进制内容同理,应在上游完成解码与分片。另外,从 EmbeddingTransform.java 的二进制处理逻辑可以推断,Embedding 侧对二进制输入已有[data, relativePath, partIndex]三列约定与分片组装缓存,属于少数内置了二进制编排的模型类 Transform,普通LLMtransform 并不具备同等能力。

如何把 Recipe 改造成自己的作业

所有 Recipe 都可以直接放进 SeaTunnel 的config目录下作为作业配置文件提交运行。仓库提供了两个可参考的作业骨架:config/v2.batch.config.template(批处理模板)与 config/v2.streaming.conf.template(流处理模板)。改造时只需把握三条主线:

  1. Source 侧:确认原始 JSON 是否保留为单列(Kafka 用format = text得到content列,NATIVE 格式则是value列;LocalFile 用file_format_type = "text")。
  2. Transform 侧:先用 JsonPath 做字段抽取与扁平化,需要向量时叠加 Embedding(有原始文本)或 DynamicCompile(已有 embedding 数组)生成FLOAT_VECTOR列。
  3. Sink 侧:JDBC 关注表结构与目标列类型匹配;Milvus 关注enable_auto_id与主键来源、向量维度是否与 collection 的 schema 一致。

生产环境建议再补充row_error_handle_way(脏数据容错)与 Milvus Sink 的data_save_modeenable_dynamic_fieldenable_upsert等参数(见 MilvusSinkOptions.java,enable_auto_id默认falsedata_save_mode默认APPEND_DATA),按目标库的写入语义做显式配置,避免默认行为与预期不符。

  • 数据集成
  • ETL
  • 大数据
  • 批处理
  • 流处理
  • 变更数据捕获

【免费下载链接】seatunnel

SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/GitHub_Trending/se/seatunnel
点击查看免费下载

相关推荐

上一篇:从零开始贡献 Azure Community-Policy:开发者必知的贡献指南
下一篇:x64dbg 调试控制命令详解:run/go/r/g 释放锁并运行程序

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询