Skip to main content

概述

Milvus Sink Connector 从 Kafka Topic 消费记录,并通过 Milvus REST API 将记录批量 upsert 到指定数据库中的一个 Collection。Kafka 记录值的字段按名称与 Collection Schema 进行匹配,匹配字段会按 Connector 支持的目标类型进行转换;Kafka 记录键、Topic、Partition、Offset 和 Header 不参与目标字段填充或 Collection 路由。 每个 Connector 实例固定写入一个 database.namecollection.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 完全同名的字段,例如显式主键 idFloatVector 字段 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 名称。topicstopics.regex 必须且只能有一个非空;所有选中 Topic 的消息值都必须兼容同一个目标 Collection。

topics.regex

按 Java 正则表达式选择要消费的 Kafka Topic。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:使用可编译且非空的 Java 正则表达式。topics.regextopics 必须且只能有一个非空;匹配到的所有 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 实现。转换结果的顶层对象必须是 Connect Structjava.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.regextopics 不能同时为非空。正则匹配到的所有 Topic 都会写入同一个 database.namecollection.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 不校验向量维度。ArrayFloat16VectorBFloat16VectorNone 没有对应转换分支;JSON 值会先被序列化为字符串,而不是作为结构化 JSON 树发送。
  • Tombstone 和发生任意字段转换异常的记录会被 Connector 自行跳过。该路径不进入 Kafka Connect 的 ErrantRecordReporter、DLQ 或 Connector 重试队列,Task 仍可能正常返回并提交越过这些记录的 Offset。
  • 一个 Worker poll 交付给 Task 的记录集合会组成一次同步 upsert 请求。Connector 没有独立的批量大小、字节大小、等待时间、拆分、异步队列或限流配置;一个被 Milvus 拒绝的实体可能导致整个请求失败。
  • HTTP 非 200、Milvus 响应码非 0、网络异常或响应解析异常都会作为不可恢复的运行时异常结束当前 Task。Connector 不区分 4295xx、认证、超时或校验错误,不处理 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.namecollection.name 是否正确,以及 Collection 是否已经处于 LoadStateLoaded。修正目标资源或权限后重启失败的 Task。

为什么 Connector 显示运行中,但 Milvus 中没有新增数据?

先确认 Collection 已关闭 auto-ID;版本 1.0.1 遇到 auto-ID Collection 时不会发起写入,也不会将批次报告为失败。随后检查消息值是否由 value.converter 转换为 StructHashMap,是否包含显式主键和 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?

topicstopics.regex 只决定 Kafka 输入范围,不参与目标路由。一个 Connector 实例中的所有记录都会写入固定的 database.namecollection.name。需要按 Topic 写入不同 Collection 时,请为不同 Topic 集合创建独立的 Connector 实例,并分别配置目标数据库和 Collection。