概述
Milvus Sink Connector 从 Kafka Topic 消费记录,并通过 Milvus REST API 将记录批量 upsert 到指定数据库中的一个 Collection。Kafka 记录值的字段按名称与 Collection Schema 进行匹配,匹配字段会按 Connector 支持的目标类型进行转换;Kafka 记录键、Topic、Partition、Offset 和 Header 不参与目标字段填充或 Collection 路由。 每个 Connector 实例固定写入一个database.name 和 collection.name。它适合把 Kafka 中已经整理为 Milvus Schema 的实体数据写入向量检索 Collection,例如同时写入业务主键、文本属性和 FloatVector 向量。目标 Collection 必须预先创建并加载,且必须关闭 auto-ID,由消息值显式提供主键。
前置条件
- 在目标 Milvus 或 Zilliz Cloud 中预先创建并加载 Collection,关闭 auto-ID,并准备可访问该数据库和 Collection、可执行 Schema 查询与实体 upsert 的令牌;消息中需要写入的字段必须与 Collection 字段同名且类型兼容,快速开始至少需要一个显式
Int64主键字段id、一个FloatVector字段vector,以及可选的标量字段。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka 和已创建并加载的 Milvus Collection,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector。<topic-name>、<milvus-rest-endpoint>、<milvus-token>、<database-name> 和 <collection-name> 替换为实际环境值。<milvus-rest-endpoint> 使用包含 http:// 或 https:// 及所需端口的基础地址;令牌形式取决于目标部署,可以是 API Token 或该部署要求的 username:password。向 Topic 写入无 Schema 的 JSON 对象时,消息值必须包含与 Collection 完全同名的字段,例如显式主键 id 和 FloatVector 字段 vector;向量元素数量必须符合 Collection 定义的维度。
配置
Milvus 连接与认证
public.endpoint
Milvus 或 Zilliz Cloud REST v2 API 的基础地址。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:Connector 不校验 URL 格式。运行时必须提供 Worker 可访问的非空 HTTP(S) 地址;Task 启动时会在该地址后拼接 Collection REST 路径并立即发起请求。
- 运行时必填:是
token
发送 Milvus REST 请求时使用的 Bearer 凭据。
- 类型:
password - 默认值:
db_admin:**** - 重要级别:高
- 有效值 / 注意事项:默认值是 ConfigDef 中的字面值,不是可用凭据。Task 在应用 ConfigDef 默认值前直接读取原始配置,因此运行时必须显式提供非空且有效的 API Token 或目标部署要求的
username:password。请使用敏感配置管理方式保存该值。 - 运行时必填:是
目标数据库与 Collection
database.name
目标 Collection 所在的 Milvus 数据库。
- 类型:
string - 默认值:
default - 重要级别:中
- 有效值 / 注意事项:使用已经存在且令牌可访问的数据库名称。Connector 不在本地校验名称或创建数据库,服务端响应决定配置是否可用。
collection.name
接收 Kafka 记录的 Milvus Collection。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:运行时必须指定已存在且处于
LoadStateLoaded状态的 Collection。Connector 不创建或加载 Collection;字段名称和类型必须与消息值兼容,并且必须关闭 auto-ID、由消息值提供主键。 - 运行时必填:是
Kafka 输入与任务
connector.class
要加载的 Milvus Sink Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:使用
com.milvus.io.kafka.MilvusSinkConnector,也可以使用 Kafka Connect 能解析到该类的别名。 - 必填:是
topics
要消费的 Kafka Topic 列表。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 有效值 / 注意事项:使用逗号分隔的 Topic 名称。
topics与topics.regex必须且只能有一个非空;所有选中 Topic 的消息值都必须兼容同一个目标 Collection。
topics.regex
按 Java 正则表达式选择要消费的 Kafka Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:使用可编译且非空的 Java 正则表达式。
topics.regex与topics必须且只能有一个非空;匹配到的所有 Topic 都写入同一个目标 Collection。
tasks.max
该 Connector 允许创建的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:必须大于或等于
1。有效并行度还受已订阅 Topic 的 Partition 数量限制;每个活跃 Task 都会独立检查同一个 Collection,并可能并发向其发起 upsert。
数据转换
key.converter
在 Connector 级别覆盖 Kafka 记录键使用的 Converter。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略或使用
null时继承 Worker 的键 Converter;显式配置时必须是可实例化的org.apache.kafka.connect.storage.Converter实现。Milvus Sink Connector 不读取转换后的记录键来填充字段或选择 Collection。
value.converter
在 Connector 级别指定 Kafka 记录值使用的 Converter。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略或使用
null时继承 Worker 的值 Converter;显式配置时必须是可实例化的org.apache.kafka.connect.storage.Converter实现。转换结果的顶层对象必须是 ConnectStruct或java.util.HashMap,字段名称和可转换值必须兼容目标 Collection Schema。Converter 自身的子配置由所选 Converter 定义。
最佳实践
按 Kafka Partition 扩展处理任务
适用业务场景:单个 Task 已成为持续写入的处理瓶颈,输入 Topic 已有多个 Partition,并且所有 Partition 的消息都遵循同一个 Collection Schema。此时可逐步增加 Task,让 Kafka Consumer Group 在 Task 之间分配 Partition。 配置示例: 在快速开始配置中增加或覆盖以下属性:3 只是示例上限,应结合输入 Partition 数量、Milvus 写入能力和实际延迟逐步调整。Task 数量超过可分配 Partition 数量不会增加记录处理并行度。多个 Task 会独立向同一 Collection 发起同步 upsert,Connector 不提供跨 Partition 或跨 Task 的全局顺序保证,也不承诺吞吐随 Task 数量线性增长。
按命名规则接入多个同构 Topic
适用业务场景:运行中的数据链路需要自动接入一组名称遵循统一规则的新 Topic,且这些 Topic 的消息字段与类型都兼容同一个目标 Collection。此时可使用正则订阅,减少逐个维护 Topic 列表的工作。 配置示例: 保留快速开始中的 Milvus 连接、Collection 和 Converter 配置,删除topics,并增加:
topics.regex 和 topics 不能同时为非空。正则匹配到的所有 Topic 都会写入同一个 database.name 和 collection.name,Connector 不会根据 Topic 名称自动选择 Collection;如果不同 Topic 需要不同 Collection 或 Schema,应拆分为不同 Connector 实例。
监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
- 版本 1.0.1 对 auto-ID 已启用的 Collection 不会发送 insert 或 upsert 请求,也不会把该批次报告为失败;因此必须关闭 auto-ID,并在消息值中显式提供主键。
- Connector 不创建或加载 Collection,也不自动迁移 Schema。每个 Task 只在启动时读取一次 Collection Schema;运行期间修改 Schema 后,需要重启 Task 才会重新读取。
- 顶层记录值只接受 Connect
Struct或具体的java.util.HashMap。顶层基本类型、JSON 字符串、List、数组、字节数组及其他 Map 实现会在 Connector 内部转换失败并被跳过。 - 字段按区分大小写的精确名称匹配。消息中不存在于 Collection Schema 的字段会被忽略,即使 Collection 启用了动态字段;Connector 不提供重命名、嵌套路径、扁平化或默认值注入。
FloatVector依赖可解析的逗号分隔列表表示,Connector 不校验向量维度。Array、Float16Vector、BFloat16Vector和None没有对应转换分支;JSON 值会先被序列化为字符串,而不是作为结构化 JSON 树发送。- Tombstone 和发生任意字段转换异常的记录会被 Connector 自行跳过。该路径不进入 Kafka Connect 的
ErrantRecordReporter、DLQ 或 Connector 重试队列,Task 仍可能正常返回并提交越过这些记录的 Offset。 - 一个 Worker poll 交付给 Task 的记录集合会组成一次同步 upsert 请求。Connector 没有独立的批量大小、字节大小、等待时间、拆分、异步队列或限流配置;一个被 Milvus 拒绝的实体可能导致整个请求失败。
- HTTP 非 200、Milvus 响应码非
0、网络异常或响应解析异常都会作为不可恢复的运行时异常结束当前 Task。Connector 不区分429、5xx、认证、超时或校验错误,不处理Retry-After,也不使用errors.tolerance或 DLQ 处理该内部 REST 失败路径。 - Kafka Offset 提交与 Milvus upsert 之间没有事务。Milvus 接受请求后、Offset 提交前发生故障时,恢复后可能重放记录;Connector 不保证恰好一次、无重复或跨 Task 全局有序。
- Connector 会在 INFO 日志中输出成功转换批次的 Collection 名称和完整数据列表,其中可能包含业务字段和大向量;应限制 Task 日志访问范围并设置合适的保留策略。
常见问题
为什么 Task 在启动时立即失败?
Task 启动时会依次检查 Collection 是否存在、是否已加载,并读取 Collection Schema。检查public.endpoint 是否为可访问的 REST 基础地址、token 是否有权限、database.name 和 collection.name 是否正确,以及 Collection 是否已经处于 LoadStateLoaded。修正目标资源或权限后重启失败的 Task。
为什么 Connector 显示运行中,但 Milvus 中没有新增数据?
先确认 Collection 已关闭 auto-ID;版本 1.0.1 遇到 auto-ID Collection 时不会发起写入,也不会将批次报告为失败。随后检查消息值是否由value.converter 转换为 Struct 或 HashMap,是否包含显式主键和 FloatVector,字段名是否与 Collection Schema 完全一致,以及日志中是否出现字段转换异常。转换失败的记录会被跳过,不会自动进入 DLQ。
为什么部分字段没有写入 Collection?
Connector 只写入启动时读取到的 Collection Schema 中存在、且名称完全匹配的字段。检查大小写、字段类型和 Task 启动后是否修改过 Schema;额外字段不会写入动态字段。若 Collection Schema 已变更,请确保消息与新 Schema 兼容,然后重启 Task 以重新读取描述信息。为什么一次 Milvus REST 错误会使 Task 进入 FAILED?
该 Connector 将非成功 HTTP 响应、非零 Milvus 响应码、网络异常和响应解析异常作为普通运行时异常抛出,而不是 Kafka Connect 的可重试异常。errors.tolerance、框架重试参数和 DLQ 不会接管这条内部 upsert 路径。检查 Task 日志和 Milvus 服务状态,修复地址、认证、Schema、请求数据或目标服务问题后再重启 Task;未提交 Offset 对应的记录可能被重新 upsert。
为什么多个 Topic 的记录都进入了同一个 Collection?
topics 和 topics.regex 只决定 Kafka 输入范围,不参与目标路由。一个 Connector 实例中的所有记录都会写入固定的 database.name 和 collection.name。需要按 Topic 写入不同 Collection 时,请为不同 Topic 集合创建独立的 Connector 实例,并分别配置目标数据库和 Collection。