概述
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。<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或 JavaMap;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。在持续运行阶段,需要为所有文档添加来源标识,并使用记录中的语言、部门等分类字段进行过滤,同时避免这些分类值重复出现在文档正文中。 配置示例:language 和 department;缺失字段会被忽略。固定的 source=kafka 会加入所有文档,记录字段则按原值写入元数据。启用排除后,language 和 department 不再生成文档内容片段,但仍保留为元数据;其他顶层非空字段继续转换为文档内容。
按 Kafka Partition 逐步扩展写入并发
适用业务场景:单 Task 已稳定写入,Kafka Topic 具有多个 Partition,积压或端到端延迟表明需要增加处理能力。此时先增加 Task 数,再谨慎调整每个 Task 的本地线程数,并观察 Vectara 限流、失败和 Offset 提交情况。 配置示例: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,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
- 非墓碑记录必须具有非空键,且值必须是顶层 Connect
Struct或 JavaMap;顶层原始类型、数组、字符串和字节不能作为文档输入。嵌套对象不会递归展开,而是以字符串形式写入文档内容。 - 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 中的最终文档状态;恢复数据时使用相同的稳定记录键,并评估重放与非原子更新对业务的影响。