> ## 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.

# Weaviate Sink Connector

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

## 概述

Weaviate Sink Connector 将 Kafka Topic 中的对象数据写入 Weaviate Collection，连接持续产生业务数据的消息链路与向量检索系统。每条非空记录的值映射为一个对象的属性，Topic 可以映射到同名 Collection，也可以集中写入指定 Collection，适用于内容索引、知识库更新和业务对象检索。

对象标识可以由客户端生成，也可以从 Kafka 消息键或记录字段派生。向量可以由目标 Collection 的向量化器生成，或从消息中的单个向量字段提取；Connector 本身不生成 Embedding。消息的序列化格式由 Kafka Connect Converter 解码，Connector 接收解码后的对象数据。

## 前置条件

* 使用 Weaviate 1.28.x 或更高版本，并提前创建目标 Collection；其属性类型和向量配置须与输入对象一致。需要服务端生成向量时，应配置相应向量化器及其模块依赖。
* Weaviate 的 HTTP 接口及配置的 gRPC 接口须对 Worker 可用；显式 gRPC 地址通常使用端口 `50051`，TLS 设置须与目标服务一致。启用认证时，准备目标服务接受的 API Key 或 OIDC 客户端凭据及写入授权。
* 使用 Connector `0.1.2` 原版发行 ZIP 时，需为该插件提供完整的独立 HTTP/2 运行时依赖；优先取得依赖完整且经过安全维护的发行包。若必须保留原版 ZIP，应将完整依赖隔离安装在同一插件位置，并在所有运行该 Connector 的 Worker 上保持一致；生产环境仍需按组织安全政策审查制品和部署方式。该限制针对该发行包，不代表所有默认部署都可直接使用。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka、输入 Topic 和已创建的 Weaviate Collection，并确认网络连通及访问权限；创建和管理操作参见[管理 Connector](../manage-connectors)。以下配置用于无需认证、HTTP 与 gRPC 均不启用 TLS 的环境，输入消息值为不带 Schema 外壳的 JSON 对象。

```properties theme={null}
connector.class=io.weaviate.connector.WeaviateSinkConnector
topics=<topic-name>
weaviate.connection.url=http://<weaviate-host>:<http-port>
weaviate.grpc.url=<weaviate-host>:<grpc-port>
collection.mapping=<collection-name>
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
```

替换 Topic、Weaviate 地址、端口和 Collection 名称后应用配置。该示例将全部输入对象写入一个指定 Collection，不使用消息键派生 ID，也不从属性中提取向量。客户端为对象生成 UUID；Collection 应配置服务端向量化器，或允许不提供向量的对象写入。若目标服务要求认证或 TLS，须按连接与认证配置调整，不能直接沿用上述无认证连接方式。

发送到 Kafka Topic 的消息值示例为以下 JSON；它是业务数据，不是 Connector 配置。目标 Collection 需接受 `title` 和 `content` 属性。

```json theme={null}
{"title":"Kafka Connect","content":"Stream object data into Weaviate."}
```

`value.converter.schemas.enable=false` 使 JsonConverter 将普通 JSON 对象转换为无 Schema 的 Map；不要发送 `schema` / `payload` 外壳，也不要以顶层字符串或数组替代对象。可在 Weaviate 中查询写入对象的属性，确认数据到达目标 Collection。默认 ID 策略在重放同一记录时会生成新 UUID，可能产生重复对象。

## 配置

### 连接

#### `weaviate.connection.url`

Weaviate HTTP 接口地址。

* **类型**：`string`
* **默认值**：`http://localhost:8080`
* **重要级别**：高
* **有效值 / 注意事项**：使用包含协议的 `http://host:port` 或 `https://host:port`，不能只填主机名。即使批量写入使用 gRPC，删除及超时后的对象读取等 HTTP 操作仍需要此地址。

#### `weaviate.grpc.url`

显式指定 Weaviate gRPC 接口地址。

* **类型**：`string`
* **默认值**：`localhost:50051`
* **重要级别**：高
* **有效值 / 注意事项**：格式为 `host:port`，不带 HTTP 协议前缀。空字符串表示不显式覆盖客户端的 gRPC 地址，不等于禁用所有 gRPC 操作。

#### `weaviate.grpc.secured`

控制显式 gRPC 连接是否启用 TLS。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：高
* **有效值 / 注意事项**：`true` 或 `false`；仅在 `weaviate.grpc.url` 非空时应用。HTTP 地址中的 `https` 不会自动将此项设为 `true`。

### 认证与请求头

#### `weaviate.auth.scheme`

选择连接 Weaviate 时使用的认证方式。

* **类型**：`string`
* **默认值**：`NONE`
* **重要级别**：高
* **有效值 / 注意事项**：使用大写 `NONE`、`API_KEY` 或 `OIDC_CLIENT_CREDENTIALS`。分别表示无认证、API Key 和 OIDC 客户端凭据认证；不要使用小写或混合大小写。

#### `weaviate.api.key`

API Key 认证使用的密钥。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **有效值 / 注意事项**：选择 `API_KEY` 时提供有效密钥。该项不是 `password` 类型，不应假定管理界面或日志自动脱敏；使用部署环境的安全凭据管理机制，不将真实密钥写入共享配置或日志。

#### `weaviate.oidc.client.secret`

OIDC 客户端凭据认证使用的 Secret。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **有效值 / 注意事项**：选择 `OIDC_CLIENT_CREDENTIALS` 时提供目标服务接受的 Secret。该项不是 `password` 类型，不应假定自动脱敏；保护其读取和导出权限。

#### `weaviate.oidc.scopes`

OIDC 客户端凭据认证请求的 Scope。

* **类型**：`list`
* **默认值**：`[openid]`
* **重要级别**：高
* **有效值 / 注意事项**：逗号分隔的 Scope，仅用于 `OIDC_CLIENT_CREDENTIALS`；按身份提供方的授权要求填写。

#### `weaviate.headers`

构建 Weaviate 客户端时附加的请求头，可用于目标向量化模块所需的认证信息。

* **类型**：`list`
* **默认值**：空列表 `[]`
* **重要级别**：中
* **有效值 / 注意事项**：逗号分隔的 `name=value`。名称和值均应非空，值中不要包含额外的 `=`，否则可能被截断；同名项后者覆盖前者。不要在共享文档或日志中暴露含凭据的请求头。

### Collection 与写入

#### `collection.mapping`

指定目标 Collection 名称或基于 Topic 的名称模板。

* **类型**：`string`
* **默认值**：`${topic}`
* **重要级别**：高
* **有效值 / 注意事项**：每个字面量 `${topic}` 替换为当前记录的 Topic 名称；固定字符串使多个 Topic 写入同一 Collection。不是逗号分隔的 Topic 到 Collection 映射表。名称不会自动清洗或调整大小写，须符合 Weaviate 命名要求，并提前创建对应 Collection。

#### `consistency.level`

传递给对象批量写入及删除操作的副本一致性级别。

* **类型**：`string`
* **默认值**：`QUORUM`
* **重要级别**：低
* **有效值 / 注意事项**：使用大写 `ALL`、`ONE` 或 `QUORUM`；所需副本及其可用性由目标 Weaviate 部署决定。该项不是 Kafka Offset 与目标写入之间的事务设置。

### 对象标识

#### `document.id.strategy`

选择从记录派生对象 ID 的策略类。

* **类型**：`class`
* **默认值**：`io.weaviate.connector.idstrategy.NoIdStrategy`
* **重要级别**：中
* **有效值 / 注意事项**：内置类为 `io.weaviate.connector.idstrategy.NoIdStrategy`、`io.weaviate.connector.idstrategy.KafkaIdStrategy` 和 `io.weaviate.connector.idstrategy.FieldIdStrategy`。`NoIdStrategy` 不提供 ID，由客户端生成随机 UUID。`KafkaIdStrategy` 将字符串消息键的 UTF-8 字节转换为名称 UUID；即使键本身看起来是 UUID，也会重新派生，非字符串键仅转换为字符串，不能假定结果是有效 UUID。`FieldIdStrategy` 将指定顶层字段转为字符串并派生名称 UUID，同时从对象属性中移除此字段。稳定 ID 应在目标 Collection 内唯一；相同键跨 Topic 汇入同一 Collection 时会产生同一 ID。自定义类须实现 `IDStrategy` 并具有可访问的无参构造器；启用删除时必须使用内置 `KafkaIdStrategy`，不接受其子类替代。

#### `document.id.field.name`

指定 `FieldIdStrategy` 读取的对象标识字段。

* **类型**：`string`
* **默认值**：`id`
* **重要级别**：中
* **有效值 / 注意事项**：仅用于 `FieldIdStrategy`，读取转换后对象的顶层属性，不支持嵌套路径。字段应始终存在且非空，并具有稳定、唯一的标量值；缺失或空值会按字符串 `null` 派生相同 ID，造成冲突。提取后该字段不再作为普通属性写入。

### 向量

#### `vector.strategy`

选择从记录提取向量的策略类。

* **类型**：`class`
* **默认值**：`io.weaviate.connector.vectorstrategy.NoVectorStrategy`
* **重要级别**：中
* **有效值 / 注意事项**：内置类为 `io.weaviate.connector.vectorstrategy.NoVectorStrategy` 和 `io.weaviate.connector.vectorstrategy.FieldVectorStrategy`。前者不提交显式向量，是否生成向量取决于 Collection；后者提取一个顶层字段作为单个向量。自定义类须实现 `VectorStrategy` 并具有可访问的无参构造器。

#### `vector.field.name`

指定 `FieldVectorStrategy` 读取的向量字段。

* **类型**：`string`
* **默认值**：`vector`
* **重要级别**：中
* **有效值 / 注意事项**：仅用于 `FieldVectorStrategy`。支持 `Float[]` 或元素为 `Float` / `Double` 的可迭代集合，`Double` 会转换为 `Float`；整数元素和标量值不能作为该向量输入。缺失或空值不提供向量，成功提取后字段从属性中移除。维度须匹配目标 Collection。ID 提取先于向量提取执行，不要与 `document.id.field.name` 使用同一字段。

### 批量处理

#### `batch.size`

控制每个 Task 的客户端批量对象数量。

* **类型**：`int`
* **默认值**：`100`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验；应结合单对象大小与目标处理能力设置。每次接收记录集合结束时也会提交剩余对象，因此实际批次可能小于此值，不是等待凑满批次才发送。

#### `pool.size`

控制每个 Task 的客户端批量处理线程池大小。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验，线程池仍需接受有效大小。此项与 `tasks.max` 独立；增加并发前评估目标负载，不应据此推导对象写入顺序或线性吞吐提升。

#### `await.termination.ms`

设置批量处理执行器关闭时的等待时长，单位为毫秒。

* **类型**：`int`
* **默认值**：`10000`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验。这是停止时的等待设置，不是单次请求、`put` 或 `flush` 的通用超时，也不是写入成功保障。

### 客户端重试参数

#### `max.connection.retries`

传递给客户端批量写入的连接错误重试次数设置。

* **类型**：`int`
* **默认值**：`3`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验。仅供客户端识别的连接错误路径使用，不代表所有 HTTP、gRPC 或单对象错误都会重试，不覆盖删除操作，也不等同于 Kafka Connect 的错误重试配置。

#### `max.timeout.retries`

传递给客户端批量写入的超时重试次数设置。

* **类型**：`int`
* **默认值**：`3`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验。仅作用于客户端识别的超时路径，不能假定任意超时都会触发重试或最终写入成功；不覆盖删除操作。

#### `retry.interval`

客户端批量重试的基础间隔，单位为毫秒。

* **类型**：`int`
* **默认值**：`2000`
* **重要级别**：低
* **有效值 / 注意事项**：未定义配置级数值范围校验。客户端对应重试路径将计数与该基础间隔相乘计算等待时间，不是指数退避；该项不会扩大可重试错误的范围。

### 删除

#### `delete.enabled`

控制是否将空消息值作为删除目标对象的请求。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：`false` 跳过空值记录；`true` 要求 `document.id.strategy=io.weaviate.connector.idstrategy.KafkaIdStrategy`。删除用相同策略从消息键派生 ID，须与原对象的写入标识一致。字符串 `"null"` 和空 JSON 对象不是 Tombstone；此项不解析 CDC 删除事件，也不保证异步写入与随后删除的完成顺序。

### Kafka Connect 运行与输入

#### `connector.class`

指定要运行的 Connector 类。

* **类型**：`string`
* **默认值**：无固定默认值，必填
* **重要级别**：高
* **有效值 / 注意事项**：使用 `io.weaviate.connector.WeaviateSinkConnector`。

#### `tasks.max`

请求运行的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`。有效消费并行度受输入分区数和分配结果约束；每个 Task 拥有独立客户端和批量处理资源，不按 Collection 自动划分 Task。

#### `topics`

指定要消费的 Topic 列表。

* **类型**：`list`
* **默认值**：空列表 `[]`
* **重要级别**：高
* **有效值 / 注意事项**：逗号分隔的 Topic 名称；与 `topics.regex` 必须且只能配置一个非空选项。若部署配置了 DLQ，不能消费该 DLQ Topic。

#### `topics.regex`

用正则表达式选择输入 Topic。

* **类型**：`string`
* **默认值**：空字符串 `""`
* **重要级别**：高
* **有效值 / 注意事项**：使用 Java 正则表达式语法，与 `topics` 互斥，不能匹配已配置的 DLQ Topic。新增匹配 Topic 时，其对象结构仍须匹配目标 Collection；使用 `${topic}` 映射时需先准备各目标 Collection。

#### `key.converter`

将 Kafka 消息键解码为 Connect 值。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未设置时使用 Worker 的 Converter 配置；显式设置须为可实例化的 Converter 类。使用 `KafkaIdStrategy` 时，解码后的键类型决定 ID 处理方式；更换 Converter 可能改变对象身份，不能仅凭 Kafka 原始字节相同判断 ID 相同。

#### `value.converter`

将 Kafka 消息值解码为 Connect 值。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未设置时使用 Worker 的 Converter 配置；显式设置须为可实例化的 Converter 类。非空值须转换为顶层 Map 或 Struct；Avro、JSON、Protobuf 等格式需要相应 Converter，不能由 Connector 直接解析。快速开始额外关闭 JsonConverter 的 Schema 外壳要求，该子设置不属于此配置的默认值。

## 最佳实践

### 扩展时按输入分区增加消费并行度

**适用业务场景**：接入已经正常运行，但多个输入分区出现持续积压，且 Weaviate 仍有处理余量时，可让多个 Task 分担消费。先确认输入至少有两个分区，并观察目标服务负载，再逐步调整并行度。

**配置示例**：在快速开始配置上添加以下设置，保留原有输入 Topic、Collection 和转换方式。

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

**关键说明**：此设置将 Task 上限从默认的一个提高到两个，分区由消费组分配，不是按 Collection 拆分。分区不足时新增 Task 无法带来额外消费并行度；每个 Task 都会增加客户端和批量处理资源。调整后同时检查分区分配、积压变化与目标资源使用情况，不同时增加 `pool.size`，便于判断变化来源。跨 Task 不提供全局对象写入顺序保证。

## 监控

### 监控内容

关注 Kafka Connect 集群健康、Connector / Task 状态、输入吞吐、消费积压及端到端延迟、Offset 提交进度和耗时、错误及重试信号，以及 Worker JVM 的堆内存、GC 和线程状态。仅在部署启用了相应错误处理时关注 DLQ 活动；不要将 Task 为 `RUNNING` 或 Offset 已提交视为目标对象全部写入成功的证明，应结合目标对象查询检查数据到达情况。

### 导入 Grafana 大盘

下载 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认监控数据已采集到 Prometheus 兼容数据源，且采集标签与大盘使用的集群、Connector、Task 和实例筛选标签一致；在 Grafana 导入 JSON 并选择对应数据源。

## 限制条件

* Connector 不创建 Collection，也不执行 Collection Schema 迁移；目标 Schema 和自动 Schema 策略决定新增属性是否被接受。
* 非空记录必须是对象，不支持顶层标量或数组；Connector 不自动拆解 CDC 外壳，也不提供字段过滤、重命名或基于字段的 Collection 路由。
* 内置向量提取仅处理单个向量，不支持命名向量或多个向量，也不在 Connector 内生成 Embedding。
* `delete.enabled=true` 只能与内置 `KafkaIdStrategy` 配合使用，不能使用 `FieldIdStrategy` 或 `NoIdStrategy`。
* 客户端以错误结果返回的批量写入或删除失败不会自动成为 Task 异常；Connector 不主动将这些失败记录交给 DLQ，不能依靠 Task 状态或框架容错设置判定这些写入成功。

## 常见问题

### Task 启动时提示认证方式或一致性级别无效

检查 `weaviate.auth.scheme` 和 `consistency.level` 的大小写。使用大写 `NONE`、`API_KEY`、`OIDC_CLIENT_CREDENTIALS` 以及 `ALL`、`ONE`、`QUORUM`；认证方式还须与目标服务匹配，并提供相应凭据。若认证已通过而连接仍失败，分别核对 HTTP 协议、gRPC 地址和 gRPC TLS 设置。

### JSON 消息无法转换或写入对象

先检查 `value.converter` 与消息格式是否一致。使用快速开始的 JsonConverter 时，消息值应是普通 JSON 对象，并设置 `value.converter.schemas.enable=false`；顶层字符串和数组不能映射为对象。再检查属性名称、类型及 Collection 的向量配置；更换 Converter 前确认已有数据格式，不通过忽略错误来替代格式修正。

### 对象写入了意料之外的 Collection

检查 `collection.mapping`。默认 `${topic}` 按原始 Topic 名称路由，不调整大小写，也不解析 `topic:collection` 形式的映射。需要集中写入时填固定 Collection 名称；需要逐 Topic 写入时，提前创建模板替换后的每个 Collection。

### 重放后出现重复对象，或不同记录对应同一个对象 ID

默认 `NoIdStrategy` 每次重新处理记录都会生成新的 UUID，因此重放可能增加对象。需要稳定身份时，先检查消息中是否已有稳定且唯一的键或标识字段，再选择相应 ID 策略并规划已有对象的清理或迁移。使用键策略时确认 Converter 输出类型，并避免多个 Topic 的相同键汇入同一 Collection；使用字段策略时检查字段是否缺失或为空。切换策略不会自动迁移原有对象的 ID，稳定 ID 也不意味着跨系统事务或写入顺序保证。

### Task 正常运行，但目标中找不到预期对象

先确认输入 Topic、消费进度及实际目标 Collection，再检查 Weaviate 的请求错误、属性类型、向量要求和认证授权。`RUNNING` 或 Offset 推进不代表目标对象已经写入；应结合 Task 和 Worker 错误信息，以及目标端的对象查询结果，确认实际写入状态。

### 发送空值后目标对象仍存在

确认发送的是 Kafka 空消息值而不是 JSON 字符串 `"null"`，并检查 `delete.enabled` 是否启用及 ID 策略是否为内置 `KafkaIdStrategy`。消息键经 Converter 转换后必须与原对象写入时的键一致；使用随机 ID 写入的对象不能用该键推导定位。还需检查目标删除请求的结果，以及同一 ID 是否仍有未完成或后续写入；不要仅凭 Task 状态判断删除已完成。
