概述
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。以下配置用于无需认证、HTTP 与 gRPC 均不启用 TLS 的环境,输入消息值为不带 Schema 外壳的 JSON 对象。title 和 content 属性。
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 和转换方式。pool.size,便于判断变化来源。跨 Task 不提供全局对象写入顺序保证。
监控
监控内容
关注 Kafka Connect 集群健康、Connector / Task 状态、输入吞吐、消费积压及端到端延迟、Offset 提交进度和耗时、错误及重试信号,以及 Worker JVM 的堆内存、GC 和线程状态。仅在部署启用了相应错误处理时关注 DLQ 活动;不要将 Task 为RUNNING 或 Offset 已提交视为目标对象全部写入成功的证明,应结合目标对象查询检查数据到达情况。
导入 Grafana 大盘
下载 AutoMQ Connect Cluster Dashboard,确认监控数据已采集到 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 状态判断删除已完成。