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

# Vectara Sink Connector

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

## 概述

Vectara Sink Connector 消费 Kafka Topic 中的记录，并将每条记录转换为 Vectara Corpus 中的文档。Kafka 记录键会转换为 Vectara 文档 ID；记录值中的顶层字段会转换为文档内容，指定字段还可以写入文档元数据。相同文档 ID 的后续记录会更新目标文档，值为 `null` 的墓碑记录会删除对应文档。

该 Connector 适合将业务内容、知识条目或其他结构化事件从 Kafka 持续写入已有的 Vectara Corpus。多个 Kafka Topic 可以写入同一个固定 Corpus，但文档 ID 不包含 Topic 或 Partition 信息，因此生产者需要提供在目标 Corpus 内稳定且不会冲突的记录键。

## 前置条件

* 提前创建目标 Vectara Corpus，并准备对该 Corpus 具有文档创建、更新和删除权限的 API Key；如果启用客户端证书认证，还需准备 Connect Worker 可读取的密钥库或其 Base64 内容。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和 Vectara Corpus，并确认网络连通和 API Key 权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=com.vectara.kafka.connect.VectaraDocumentSinkConnector
topics=vectara-documents
tasks.max=1
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
api.key=<vectara-api-key>
customer.id=<vectara-customer-id>
corpus.id.location=Config
corpus.id=<vectara-corpus-id>
```

将 `<vectara-api-key>`、`<vectara-customer-id>` 和 `<vectara-corpus-id>` 替换为实际值。向 `vectara-documents` 写入记录时使用非空字符串键，并让记录值是顶层 JSON 对象；例如键 `article-1001` 会成为目标文档 ID。`value.converter.schemas.enable=false` 是 JSON Converter 的配置，用于把普通 JSON 对象转换为 Connector 支持的顶层 `Map`。

## 配置

### Vectara 连接与账号

#### `api.key`

用于调用 Vectara API 的凭证。

* **类型**：`password`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。配置解析允许空值，但空值无法形成有效认证。请通过受控的 Kafka ConfigProvider 管理该凭证，不要写入日志或公开配置仓库。
* **必填**：是

#### `api.url`

Vectara API 的基础地址。

* **类型**：`string`
* **默认值**：`https://api.vectara.io/`
* **重要级别**：高
* **有效值 / 注意事项**：必须是绝对 URI。配置校验不限制协议、不检查主机，也不验证地址是否可访问；生产环境应使用受信任的 HTTPS 端点。

#### `customer.id`

Vectara 客户 ID。

* **类型**：`long`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填，必须能解析为 `long`。1.0.4 会解析并保存该值，但没有确认该值参与 API 请求、Header 或 Corpus 路由；不要依赖它实现账号路由。
* **必填**：是

### Corpus 路由

#### `corpus.id.location`

指定从 Connector 配置还是记录字段取得目标 Corpus ID。

* **类型**：`string`
* **默认值**：`Config`
* **重要级别**：高
* **有效值 / 注意事项**：仅接受区分大小写的 `Config` 或 `Field`。1.0.4 应使用 `Config`，并通过 `corpus.id` 指定固定 Corpus；该版本的 `Field` 路由不会按 `corpus.field` 读取字段，且无法处理墓碑删除。

#### `corpus.id`

固定目标 Corpus ID。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：当 `corpus.id.location=Config` 时必填且不能为空。所有订阅 Topic 的记录都会写入该 Corpus，并共享同一个文档 ID 命名空间。
* **必填**：条件必填

#### `corpus.field`

声明记录值中用于取得 Corpus ID 的顶层字段名。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：1.0.4 虽然公开此配置，但运行时不会读取其值；配置它不会改变 `Field` 模式实际查找的字段。请改用 `corpus.id.location=Config` 和非空 `corpus.id`。

### 文档与元数据

#### `document.metadata.fields`

从记录值的顶层字段中选择要写入 Vectara 文档元数据的字段。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：高
* **有效值 / 注意事项**：使用逗号分隔、区分大小写的顶层字段名。Map 中缺失的字段会被忽略；Struct 只检查其 Schema 中存在的字段。重复名称会去重。记录字段与 `metadata.default.*` 产生同名元数据时，记录字段值优先。

#### `exclude.metadata.fields.from.document`

是否从文档正文中排除 `document.metadata.fields` 选中的字段。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：`false` 表示选中字段同时出现在元数据和文档内容中；`true` 表示这些字段只保留为元数据。该配置不会排除仅由 `metadata.default.*` 添加的元数据键。

#### `metadata.default.*`

为每个 Vectara 文档添加固定元数据。将 `*` 替换为实际元数据键，例如 `metadata.default.source=kafka` 会生成元数据 `source=kafka`。

* **类型**：原始配置值
* **默认值**：无
* **重要级别**：不适用
* **有效值 / 注意事项**：这是动态前缀配置，不属于 Connector 的固定 ConfigDef 项，因此没有统一类型、默认值或校验规则。前缀会被移除；同名的 `document.metadata.fields` 记录值会覆盖固定值。不要通过该前缀写入敏感信息。

### 批处理与并发

#### `document.batch.timeout.seconds`

每次处理一批 Kafka 记录时，等待该批异步文档任务的最长时间，单位为秒。

* **类型**：`long`
* **默认值**：`60`
* **重要级别**：低
* **有效值 / 注意事项**：配置没有最小值校验，但应使用正数。该值不是 HTTP 请求超时、任务取消期限或 Offset 提交屏障；超时后仍未完成的远程操作可能继续执行。

#### `callback.executor.pool.size`

每个 Sink Task 用于执行 Vectara 请求的本地线程数。

* **类型**：`int`
* **默认值**：`10`
* **重要级别**：低
* **有效值 / 注意事项**：应大于或等于 `1`；`0` 或负数虽然能通过配置类型解析，但会导致 Task 启动失败。总并发还会随活动 Task 数量增长，且同一批记录可能并发完成。

#### `max.requests`

声明允许的最大请求数。

* **类型**：`int`
* **默认值**：`10`
* **重要级别**：低
* **有效值 / 注意事项**：有效范围为 `1` 到 `20`，含边界。1.0.4 会校验并保存该值，但没有确认它会限制本地执行器或客户端请求并发；调整吞吐时应以 `tasks.max` 和 `callback.executor.pool.size` 为主要控制项。

### TLS 客户端密钥库

#### `ssl.keystore.location`

指定客户端密钥库的来源。

* **类型**：`string`
* **默认值**：`None`
* **重要级别**：高
* **有效值 / 注意事项**：仅接受区分大小写的 `File`、`Inline` 或 `None`。`None` 不加载客户端密钥库；`File` 使用 `ssl.keystore.path`；`Inline` 使用 `ssl.key.inline`。密钥库类型使用 JVM 默认值。

#### `ssl.keystore.path`

客户端密钥库文件路径。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：当 `ssl.keystore.location=File` 时条件必填，文件必须可由 Connect Worker 进程读取，并与 JVM 默认 KeyStore 类型及配置密码匹配。
* **必填**：条件必填

#### `ssl.keystore.password`

加载 File 或 Inline 客户端密钥库时使用的密码。

* **类型**：`password`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：当密钥库使用非空密码时条件必填。`ssl.keystore.location=None` 时忽略。请通过受控的 Kafka ConfigProvider 管理该凭证。
* **必填**：条件必填

#### `ssl.key.inline`

Base64 编码的客户端密钥库内容。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：当 `ssl.keystore.location=Inline` 时条件必填。内容必须可解码为 JVM 默认 KeyStore 类型的密钥库。虽然配置类型是字符串，但它包含私钥材料，必须按敏感信息管理。
* **必填**：条件必填

### 输入订阅

#### `topics`

要消费的 Kafka Topic 列表。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：高
* **有效值 / 注意事项**：与 `topics.regex` 二选一。使用逗号分隔的非空 Topic 名称。多个 Topic 写入同一 Corpus 时必须协调记录键，避免不同 Topic 产生相同文档 ID。
* **必填**：条件必填

#### `topics.regex`

按 Java 正则表达式订阅 Kafka Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：与 `topics` 二选一，必须是合法且非空的 Java 正则表达式。配置 DLQ 时，该表达式不能匹配 DLQ Topic。
* **必填**：条件必填

### Connector 身份与任务

#### `connector.class`

要加载的 Vectara Sink Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：1.0.4 必须使用 `com.vectara.kafka.connect.VectaraDocumentSinkConnector`。
* **必填**：是

#### `tasks.max`

该 Connector 最多可以创建的 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须大于或等于 `1`。实际活动 Task 数不会超过已订阅 Kafka Partition 数；每个 Task 还会创建 `callback.executor.pool.size` 个本地执行线程。

#### `tasks.max.enforce`

是否启用 Kafka Connect 对 `tasks.max` 的框架约束。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：Kafka Connect 3.9.1 中仍可配置，但已标记为弃用并计划在后续主版本移除。保持启用并通过 `tasks.max` 控制任务上限。
* **已弃用**：是

### 数据转换

#### `key.converter`

指定 Kafka 记录键的 Converter。转换后的键通过 `toString()` 生成 Vectara 文档 ID。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的键 Converter；配置时必须是可实例化的具体 Converter 类。每条写入和墓碑删除记录都必须产生非空、稳定且在目标 Corpus 内唯一的键字符串。

#### `value.converter`

指定 Kafka 记录值的 Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的值 Converter；配置时必须是可实例化的具体 Converter 类。非墓碑记录必须转换为顶层 Connect `Struct` 或 Java `Map`；Map 的字段键必须是字符串。顶层字符串、数字、数组和字节不受支持。

### 错误处理

#### `errors.tolerance`

指定 Kafka Connect 对可容忍记录错误的处理策略。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：仅接受 `none` 或 `all`。该配置不能保证所有 Vectara 异步写入或删除失败都可重试、跳过或进入 DLQ；仍需结合 Task 日志和目标端结果确认投递状态。

#### `errors.deadletterqueue.topic.name`

用于保存符合 Kafka Connect DLQ 条件的错误记录的 Topic 名称。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：空字符串表示不启用 DLQ。非空 DLQ Topic 不能同时出现在 `topics` 中，也不能被 `topics.regex` 匹配。配置 DLQ 不代表 Connector 自身捕获的所有异步失败都会发布到该 Topic。

#### `errors.deadletterqueue.topic.replication.factor`

Kafka Connect 创建缺失的 DLQ Topic 时使用的副本因子。

* **类型**：`short`
* **默认值**：`3`
* **重要级别**：中
* **有效值 / 注意事项**：仅在已配置的 DLQ Topic 不存在并由 Kafka Connect 创建时使用；该值必须适合 Kafka 集群可用 Broker 数量和副本策略。

## 最佳实践

### 将业务分类作为元数据并从正文排除

**适用业务场景**：Connector 已经能够将内容写入固定 Corpus。在持续运行阶段，需要为所有文档添加来源标识，并使用记录中的语言、部门等分类字段进行过滤，同时避免这些分类值重复出现在文档正文中。

**配置示例**：

```properties theme={null}
connector.class=com.vectara.kafka.connect.VectaraDocumentSinkConnector
topics=vectara-documents
tasks.max=1
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
api.key=<vectara-api-key>
customer.id=<vectara-customer-id>
corpus.id.location=Config
corpus.id=<vectara-corpus-id>
metadata.default.source=kafka
document.metadata.fields=language,department
exclude.metadata.fields.from.document=true
```

**关键说明**：每条 JSON 对象应在顶层提供 `language` 和 `department`；缺失字段会被忽略。固定的 `source=kafka` 会加入所有文档，记录字段则按原值写入元数据。启用排除后，`language` 和 `department` 不再生成文档内容片段，但仍保留为元数据；其他顶层非空字段继续转换为文档内容。

### 按 Kafka Partition 逐步扩展写入并发

**适用业务场景**：单 Task 已稳定写入，Kafka Topic 具有多个 Partition，积压或端到端延迟表明需要增加处理能力。此时先增加 Task 数，再谨慎调整每个 Task 的本地线程数，并观察 Vectara 限流、失败和 Offset 提交情况。

**配置示例**：

```properties theme={null}
connector.class=com.vectara.kafka.connect.VectaraDocumentSinkConnector
topics=vectara-documents
tasks.max=2
callback.executor.pool.size=4
document.batch.timeout.seconds=120
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
api.key=<vectara-api-key>
customer.id=<vectara-customer-id>
corpus.id.location=Config
corpus.id=<vectara-corpus-id>
```

**关键说明**：有效 Task 数受 Kafka Partition 数限制，总的本地执行线程上限约为活动 Task 数乘以 `callback.executor.pool.size`。从较小值开始逐步增加，并根据目标端限流、错误率和积压调整。`document.batch.timeout.seconds` 只控制 `put` 对异步任务的等待时间，不会取消超时后的请求，也不能保证远程写入完成后才提交 Offset。对同一文档 ID 的连续更新可能并发或重排，因此不要依赖处理完成顺序表达业务状态。

## 监控

### 监控内容

关注 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 并选择对应数据源。

## 限制条件

* 非墓碑记录必须具有非空键，且值必须是顶层 Connect `Struct` 或 Java `Map`；顶层原始类型、数组、字符串和字节不能作为文档输入。嵌套对象不会递归展开，而是以字符串形式写入文档内容。
* 1.0.4 的 `corpus.field` 不参与运行时字段选择；`corpus.id.location=Field` 实际查找名为 `Field` 的顶层字段，并且墓碑记录没有可用于路由的值。生产配置应使用固定的 `Config` 路由。
* 文档 ID 仅来自记录键的字符串表示，不包含 Kafka Topic 或 Partition。多个 Topic 写入同一 Corpus 时，相同键字符串会指向同一个文档。
* 文档更新采用删除后创建的非原子流程；同一文档 ID 的并发更新或删除可能重排，并可能短暂不存在。Connector 不提供端到端事务或恰好一次投递保证。
* `put` 等待超时后，远程操作仍可能继续，且 `flush` 不等待未完成请求；Kafka Offset 可能在对应 Vectara 操作完成或失败前提交。进程故障也可能导致已成功的记录在恢复后重放。
* `max.requests` 和必填的 `customer.id` 在 1.0.4 中会被解析，但没有确认会分别限制请求并发或参与 Vectara API 请求。不要依赖这两个配置实现相应运行时控制。

## 常见问题

### 为什么 Task 报告记录值类型不受支持？

检查 `value.converter` 的输出，而不只是 Kafka 中的原始序列化格式。非墓碑记录必须转换为顶层 `Map` 或 `Struct`；使用普通 JSON 对象时，可按快速开始配置 `JsonConverter` 并设置 `value.converter.schemas.enable=false`。还需确认 `Map` 的顶层字段键都是字符串，且记录键不为空。

### 为什么配置 `corpus.field` 后记录没有写入预期 Corpus？

这是 1.0.4 的字段路由边界：该版本不会读取 `corpus.field` 的值，而会在 `Field` 模式下查找名为 `Field` 的顶层字段。将 `corpus.id.location` 改为 `Config`，并把目标 Corpus ID 写入非空的 `corpus.id`；需要多个 Corpus 时，为不同输入范围创建分别使用固定 Corpus 的 Connector 实例。

### 为什么元数据字段仍然出现在文档正文中？

`document.metadata.fields` 只负责选择元数据，默认不会从正文移除这些字段。将 `exclude.metadata.fields.from.document=true`，并确认字段名与记录顶层字段完全一致且大小写相同。`metadata.default.*` 添加的固定元数据不会自动对应或排除正文中的同名字段。

### 为什么提高 `max.requests` 后吞吐没有明显变化？

1.0.4 会校验 `max.requests` 的范围，但没有确认它会控制执行器或客户端并发。优先检查 Kafka Partition 数、活动 Task 数和积压，再通过 `tasks.max` 与 `callback.executor.pool.size` 逐步调整并发；同时观察 Vectara 限流、错误率和 Worker 资源使用情况。

### 为什么 Kafka Offset 已推进，但 Vectara 文档仍缺失或稍后才出现？

Connector 会异步执行远程操作，批次等待超时不会取消尚未完成的请求，`flush` 也不会等待这些请求。因此 Offset 与目标端完成状态之间可能存在窗口。检查 Connector 和 Task 状态、相关时间段日志、错误与 DLQ 活动以及 Vectara 中的最终文档状态；恢复数据时使用相同的稳定记录键，并评估重放与非原子更新对业务的影响。
