> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Milvus Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Milvus Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.milvus.io.kafka.MilvusSinkConnector
topics=<topic-name>
public.endpoint=<milvus-rest-endpoint>
token=<milvus-token>
database.name=<database-name>
collection.name=<collection-name>
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
```

将 `<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` 实现。转换结果的顶层对象必须是 Connect `Struct` 或 `java.util.HashMap`，字段名称和可转换值必须兼容目标 Collection Schema。Converter 自身的子配置由所选 Converter 定义。

## 最佳实践

### 按 Kafka Partition 扩展处理任务

**适用业务场景**：单个 Task 已成为持续写入的处理瓶颈，输入 Topic 已有多个 Partition，并且所有 Partition 的消息都遵循同一个 Collection Schema。此时可逐步增加 Task，让 Kafka Consumer Group 在 Task 之间分配 Partition。

**配置示例**：

在快速开始配置中增加或覆盖以下属性：

```properties theme={null}
tasks.max=3
```

**关键说明**：`3` 只是示例上限，应结合输入 Partition 数量、Milvus 写入能力和实际延迟逐步调整。Task 数量超过可分配 Partition 数量不会增加记录处理并行度。多个 Task 会独立向同一 Collection 发起同步 upsert，Connector 不提供跨 Partition 或跨 Task 的全局顺序保证，也不承诺吞吐随 Task 数量线性增长。

### 按命名规则接入多个同构 Topic

**适用业务场景**：运行中的数据链路需要自动接入一组名称遵循统一规则的新 Topic，且这些 Topic 的消息字段与类型都兼容同一个目标 Collection。此时可使用正则订阅，减少逐个维护 Topic 列表的工作。

**配置示例**：

保留快速开始中的 Milvus 连接、Collection 和 Converter 配置，删除 `topics`，并增加：

```properties theme={null}
topics.regex=milvus-events-.*
```

**关键说明**：`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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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。
