☰
SeaTunnel Milvus Sink Connector 实战指南:向 Milvus 与 Zilliz Cloud 写入向量数据
2026/9/28 2:49:23 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

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

本文以官方文档 docs/en/connector-v2/sink/Mivlus.md 为核心骨架,并结合 SeaTunnel 仓库中connector-milvus模块的源码实现,系统讲解 Milvus Sink Connector 的能力边界、参数语义、类型映射与写入原理。读完本文,你将能够独立完成一个「从任意数据源 → SeaTunnel → Milvus/Zilliz Cloud 向量库」的批式同步任务,并理解其批量写入、upsert 与 exactly-once 等机制在源码层面是如何落地的。

一、连接器概述与功能特性

Milvus Sink Connector 用于将 SeaTunnel 作业中的数据写入Milvus或Zilliz Cloud(两者兼容同一套 RESTful/gRPC 访问协议,均通过url与token建立连接)。连接器的工厂标识(factoryIdentifier)为Milvus,对应源码 MilvusSinkFactory.java,因此任务配置中的 sink 名即为Milvus。

依据官方文档的 Connector V2 特性说明,该连接器支持的能力矩阵如下:

特性支持情况说明
batch(批式)✅ 支持数据有界,作业处理完即结束,适合离线/批式向量入库场景
exactly-once(精确一次)✅ 支持通过 sink 的 committer 与状态恢复机制实现,详见下文第六节
column projection(列投影)❌ 不支持无法仅写入指定列;如需裁剪字段,应在 sink 前通过 Transform V2 的字段映射等转换 完成

二、写入链路与源码架构

从源码结构看,Milvus Sink 的写入链路由四层组成,对应 connector-milvus 模块 下的sink包:

  1. MilvusSink:连接器的顶层实现,实现SeaTunnelSink与SupportSaveMode接口(MilvusSink.java)。它负责创建 Writer、恢复 Writer、创建 Committer,并基于schema_save_mode/data_save_mode通过MilvusCatalog生成SaveModeHandler来管理建库建表。
  2. MilvusSinkWriter:每个并行子任务对应一个 Writer,持有 Milvus 客户端连接(MilvusSinkWriter.java)。它使用 Milvus Java SDK V2 的ConnectConfig以uri+token建连,并驱动内部的批量写入器。
  3. MilvusBufferBatchWriter:缓冲批量写入器(MilvusBufferBatchWriter.java),负责把SeaTunnelRow转换为 Milvus 的 JSON 数据、按batch_size攒批、并执行insert/upsert请求。
  4. MilvusSinkCommitter:负责两阶段提交中的 commit 阶段,配合 Writer 的prepareCommit实现精确一次语义。

此外,catalog包下的 MilvusCatalog.java 承担「表不存在时自动建表」的能力:createTableInternal会把 SeaTunnel 的CatalogTable列定义转换为 Milvus 的FieldType,并通过createIndexInternal为向量字段创建索引(CreateIndexParam中指定index_type与metric_type);建集合时默认使用ConsistencyLevelEnum.BOUNDED一致性级别,并依据配置决定是否开启 dynamic field。

三、数据类型映射

文档给出了 Milvus 与 SeaTunnel 之间的类型映射表,完整继承如下:

Milvus Data TypeSeaTunnel Data Type
INT8TINYINT
INT16SMALLINT
INT32INT
INT64BIGINT
FLOATFLOAT
DOUBLEDOUBLE
BOOLBOOLEAN
JSONSTRING
ARRAYARRAY
VARCHARSTRING
FLOAT_VECTORFLOAT_VECTOR
BINARY_VECTORBINARY_VECTOR
FLOAT16_VECTORFLOAT16_VECTOR
BFLOAT16_VECTORBFLOAT16_VECTOR
SPARSE_FLOAT_VECTORSPARSE_FLOAT_VECTOR

该映射在源码中有两套对应实现,分别服务于「建表」与「写数据」两个方向:

3.1 建表方向(SeaTunnel → Milvus)

MilvusConvertUtils.convertSqlTypeToDataType 负责将 SeaTunnel 的 SQL 类型转为 Milvus 的DataType。需要注意几个在源码中体现的细节:

  • STRING→VarChar;DATE、ROW也会被转成VarChar(其中 DATE 字段最大长度固定 20,ROW 固定 65535,见 MilvusCatalog.java);
  • MAP在建表时被转成 Milvus 的JSON类型(MilvusCatalog.java);
  • 字符串字段:未声明长度时默认max_length = 512,声明了长度则取columnLength / 4(因为 SeaTunnel 按字符数 ×4 换算字节,见 MilvusCatalog.java);
  • 数组字段:需指定元素类型,max_capacity固定为 4095(MilvusCatalog.java);
  • 向量字段(FLOAT_VECTOR / BINARY_VECTOR / FLOAT16_VECTOR / BFLOAT16_VECTOR):维度取自列的scale(dim),建表时必须正确声明;
  • 主键字段通过withPrimaryKey(true)标记,autoID由主键定义或enable_auto_id配置决定(MilvusCatalog.java)。

3.2 写数据方向(Milvus 行 → JSON)

convertBySeaTunnelType 负责把SeaTunnelRow中的每个字段转换为 Milvus SDK 可识别的 Java 对象:

  • 数值/布尔:TINYINT/SMALLINT/INT/BIGINT/FLOAT/DOUBLE/BOOLEAN分别解析为对应的 Java 基本类型;
  • STRING与DATE:直接toString;
  • FLOAT_VECTOR:由Object[]转为List<Float>;
  • ARRAY:根据元素类型(STRING/INT/BIGINT/FLOAT/DOUBLE)转为Arrays.asList(...);
  • ROW:序列化为 JSON 字符串;MAP:序列化为 JSON 字符串。

对于不在此范围内的类型,写入时会抛出NOT_SUPPORT_TYPE异常,建表时也会抛出CatalogException提示Not support convert to milvus type。

四、Sink 参数详解

文档给出的全部 Sink 参数如下表,其中标注了默认值与是否必填:

NameTypeRequiredDefaultDescription
urlStringYes-连接 Milvus 或 Zilliz Cloud 的地址
tokenStringYes-认证信息,格式为User:password
databaseStringNo-写入的目标数据库;不配置时使用数据源所属数据库
schema_save_modeenumNoCREATE_SCHEMA_WHEN_NOT_EXIST表不存在时自动建表
enable_auto_idbooleanNofalse主键列是否启用 autoId
enable_upsertbooleanNofalse使用 upsert 而非 insert 写入
enable_dynamic_fieldbooleanNotrue建表时是否启用 dynamic field
batch_sizeintNo1000每次批量写入的行数

以上参数在 MilvusSinkConfig.java 中有完整定义。结合源码,可以补充以下几点文档未展开的实现细节:

  • url / token 为必填项:在 MilvusSinkFactory.optionRule 中通过required(...)强制校验,缺少任一参数作业会直接校验失败;两者同时被用于 Writer 建连(ConnectConfig)与 Catalog 建连(ConnectParam)。
  • database 的缺省行为:未配置时,MilvusSinkFactory.renameCatalogTable 会沿用源表所属数据库名,因此写入目标库名与源库名一致。
  • data_save_mode(源码新增参数):文档表格中未列出,但当前仓库源码已支持该参数,默认APPEND_DATA,可选DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS(MilvusSinkConfig.java),用于定义已有数据时的处理策略,与schema_save_mode共同构成 Save Mode 体系。
  • enable_upsert 的默认值差异提示:文档表格标注默认false,而当前仓库源码 MilvusSinkConfig.java 中实际默认值为true。以仓库源码为准时,未显式配置enable_upsert的作业会走 upsert 路径。升级或迁移版本时建议显式声明该参数,避免行为漂移。
  • enable_auto_id 的优先级:若表的主键(PrimaryKey)中已声明enableAutoId,则以主键声明为准;否则回落到该配置项(MilvusSinkWriter.getAutoId)。
  • enable_dynamic_field:控制建表时是否启用 Milvus 的 dynamic field(MilvusCatalog.java),开启后可写入 schema 之外的动态字段。
  • batch_size:决定单次insert/upsert请求的行数上限,也是 Writer 内存缓冲的容量(new ArrayList<>(batchSize)),攒满即触发 flush(MilvusBufferBatchWriter.java)。

五、任务配置示例

5.1 文档最小示例

文档提供的sink配置如下(可直接复制使用):

sink { Milvus { url = "http://127.0.0.1:19530" token = "username:password" batch_size = 1000 } }

5.2 完整作业示例(带 Source)

将最小示例补全为一个端到端可运行的批式作业(以 Fake 数据源为例,便于本地验证;向量字段由上游 Transform 生成):

env { execution.parallelism = 2 job.mode = "BATCH" } source { Fake { schema = { fields { id = BIGINT content = STRING embedding = "FLOAT_VECTOR" } } rows = [ { fields = [1, "hello seatunnel", [0.1, 0.2, 0.3, 0.4]] } { fields = [2, "vector database", [0.5, 0.6, 0.7, 0.8]] } ] } } transform { # 如需对字段做裁剪、重命名或类型转换,可在此处使用 Transform V2 # 参见 docs/en/transform-v2/field-mapper.md } sink { Milvus { url = "http://127.0.0.1:19530" token = "username:password" database = "default" schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode = "APPEND_DATA" enable_auto_id = false enable_upsert = true enable_dynamic_field = true batch_size = 1000 } }

运行方式:将上述配置保存为milvus.conf,使用项目根目录的启动脚本执行:

# 使用已构建好的发行包(seatunnel-core/seatunnel-starter) bin/start-seatunnel.sh --config milvus.conf # 本地开发环境亦可通过 maven wrapper 运行 starter 模块验证 ./mvnw -pl seatunnel-core/seatunnel-starter -am package

提示:embedding列需要声明为 SeaTunnel 的向量类型FLOAT_VECTOR,建表时其维度由列定义(scale)决定,请确保上游数据维度与表定义一致。

六、批量写入、upsert 与一致性机制

6.1 攒批与 flush 时机

MilvusSinkWriter.write 每收到一行数据就将其交给MilvusBufferBatchWriter.addToBatch缓存;当writeCount >= batchSize时触发flush()。flush 是synchronized的,会把缓存中的 JSON 列表一次性提交,随后清空缓存(MilvusBufferBatchWriter.java)。此外,prepareCommit()与close()时也会强制 flush,确保 checkpoint 与作业结束时缓存不残留。

6.2 insert 与 upsert 的选择逻辑

writeData2Collection 中,只有当enableUpsert && !autoId时才使用UpsertReq(upsert 依赖用户提供主键),否则使用InsertReq。也就是说:

  • 主键开启 autoId(自动生成)时,即使enable_upsert = true也会退化为 insert,此时主键字段无需随数据写入(buildMilvusData 会跳过主键字段);
  • 字段值为 null 时写入会直接抛FIELD_IS_NULL异常,属于默认的严格校验行为(MilvusBufferBatchWriter.java)。

6.3 exactly-once 的实现路径

文档声明该 Sink 支持 exactly-once。从源码看,MilvusSink 实现了两阶段提交所需的完整接口:

  • createWriter/restoreWriter:支持从MilvusSinkState状态列表恢复 Writer(第 69-72 行);
  • getWriterStateSerializer/getCommitInfoSerializer:状态与提交信息均可序列化,供 checkpoint 持久化(第 75-87 行);
  • createCommitter:返回MilvusSinkCommitter执行 commit(第 80-82 行);
  • Writer 侧prepareCommit()在 checkpoint 前强制 flush(MilvusSinkWriter.java)。

即:Writer 在 checkpoint 前完成批量提交并记录状态,任务重启后通过状态恢复 Writer,再由 Committer 完成最终提交,从而保证每条数据只被写入一次。

七、相关源码索引

便于继续深入研读的仓库文件:

  • 连接器配置定义:MilvusSinkConfig.java
  • Sink 工厂与参数校验:MilvusSinkFactory.java
  • Sink 主实现(两阶段提交、Save Mode):MilvusSink.java
  • Writer 与批量缓冲写入器:MilvusSinkWriter.java、MilvusBufferBatchWriter.java
  • 类型转换工具:MilvusConvertUtils.java
  • 建表/建索引 Catalog 实现:MilvusCatalog.java
  • Connector V2 特性定义:connector-v2-features.md

同一模块下还包含 Milvus Source Connector(MilvusSource.java 等),可实现「Milvus → Milvus」或其他 Milvus 参与的同步链路;关于变换(Transform)能力,可参考 Transform V2 文档目录 与 SQL 变换。

  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

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

相关推荐

上一篇:FS-Blog错误处理指南:Spring Boot全局异常处理与Log4j2日志系统配置详解
下一篇:Enquirer与Zod集成:类型安全的输入验证实现

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

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

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

立即咨询