Skip to main content

概述

ArangoDB Sink Connector 消费 Kafka Topic 中的记录,并把记录写入一个指定的 ArangoDB 数据库。每条记录根据 Kafka Topic 名称路由到 ArangoDB Collection:Topic 名称不含句点时使用完整名称,包含句点时使用最后一个句点后的后缀。例如,dbserver1.inventory.customers 会写入 customers Collection。 Connector 使用 Kafka 记录 Key 的第一个字段作为 ArangoDB 文档 _key。非 null 的对象 Value 会插入新文档或完整替换同一 _key 的现有文档;普通模式下,Value 为 null 的墓碑记录会按 _key 删除文档。它适合把业务对象、缓存视图或变更事件持续同步到 ArangoDB,但不提供字段级局部更新、关系边映射或跨 Kafka 与 ArangoDB 的事务。

前置条件

  • 使用 ArangoDB 3.4 或更高版本,预先创建目标 Database,并为每个被订阅 Topic 预先创建按 Topic 最后一个句点后缀命名的 Collection;Connector 不会自动创建这些资源,所用账号需要对目标 Database 和 Collection 具有文档写入与删除权限。

授权许可

使用 MIT License。

快速开始

提前准备 Connect Cluster、Kafka、目标 ArangoDB Database 和 customers Collection,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
替换 ArangoDB 地址、端口、账号、密码和 Database 占位符。向 customers Topic 写入对象 Key,例如 {"id":"customer-1001"},并写入对象 Value,例如 {"id":"customer-1001","name":"Alice","tier":"gold"};Connector 使用 Key 中第一个字段的值作为 _key,并从写入文档中移除 Value 内同名的 id 字段。后续发送同一 Key 的新对象 Value 会完整替换原文档,而不是只更新发生变化的字段。生产环境部署前还应评估本版本记录原始配置时可能暴露密码的风险。

配置

Connector 身份与 Task

connector.class

要加载的 ArangoDB Sink Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 io.github.jaredpetersen.kafkaconnectarangodb.sink.ArangoDbSinkConnector。不要使用安装包快速开始文件中的旧类名。
  • 必填:是

tasks.max

Connector 请求创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须至少为 1。实际并行度受 Topic 分区分配限制;多个 Task 之间没有全局顺序或同一文档键的冲突协调。

Topic 订阅

topics

要消费的 Kafka Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔的字面 Topic 名称。必须在 topicstopics.regex 中恰好选择一个非空配置;每个 Topic 的最后一个句点后缀必须对应已存在的 ArangoDB Collection。
  • 必填conditional(未配置 topics.regex 时必填)

topics.regex

用于选择 Kafka Topic 的正则表达式。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项:使用 Java 正则表达式语法。必须在 topicstopics.regex 中恰好选择一个非空配置;不要让表达式匹配错误处理使用的 DLQ Topic。
  • 必填conditional(未配置 topics 时必填)

ArangoDB 连接与目标

arangodb.host

ArangoDB 服务端主机名或地址。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:填写 Connect Worker 可访问的 ArangoDB 主机名或地址。此配置没有非空或可达性校验。
  • 必填:是

arangodb.port

ArangoDB 服务端端口。
  • 类型int
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:填写实际监听端口。Connector 只校验值能否解析为整数,不校验正数或 TCP 端口范围。
  • 必填:是

arangodb.database.name

所有记录写入的 ArangoDB Database 名称。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:Database 必须预先存在,且所有订阅 Topic 共用这一个 Database。Collection 名称由 Topic 后缀决定,不能在此配置映射。
  • 必填:是

ArangoDB 身份认证

arangodb.user

连接 ArangoDB 使用的用户名。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:账号需要具备目标 Database 和 Collection 的文档写入与删除权限。此配置没有非空或认证预检。
  • 必填:是

arangodb.password

连接 ArangoDB 使用的密码。
  • 类型password
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项:省略时使用空字符串。1.0.4 会在 INFO 日志中记录原始配置 Map,可能暴露这里提供的密码;限制 Worker 日志访问和保留范围,并在生产部署前采用已修复的构建或其他经评估的缓解措施。

记录转换

key.converter

Connector 级 Kafka 记录 Key Converter 类。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 的 Key Converter。Converter 必须生成至少包含一个字段且首字段值非 null 的对象 Key;Connector 只使用对象的第一个字段值作为 ArangoDB _key,不支持原始字符串、数值、字节数组或稳定的复合键映射。

value.converter

Connector 级 Kafka 记录 Value Converter 类。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 的 Value Converter。转换后的非 null Value 必须是一个对象;数组和标量不能作为 ArangoDB 文档写入。普通模式下,转换后的 null Value 表示删除。

最佳实践

扩展接入时按 Topic 规则写入多个 Collection

适用业务场景:已经在一个 ArangoDB Database 中按业务对象建立多个 Collection,并希望后续符合统一命名规则的 Topic 由同一个 Connector 订阅。以下配置使用 inventory. 前缀匹配 Topic;每个 Topic 仍按最后一个句点后的名称写入对应 Collection。
关键说明:移除快速开始中的 topics,因为它与 topics.regex 互斥。启用前先为每个可能匹配的 Topic 后缀创建同名 Collection,例如 inventory.customersinventory.orders 分别需要 customersorders;同时确认表达式不会意外订阅内部 Topic 或 DLQ。所有匹配 Topic 都写入同一个 arangodb.database.name,无法按 Topic 选择不同 Database。

接入 CDC 事件时统一处理更新和删除

适用业务场景:上游 CDC 系统输出带 beforeafter 的对象,希望把 after 作为当前完整文档写入 ArangoDB,并在 after=null 时删除对应文档。以下配置在快速开始基础上增加 Connector 自带的 CDC SMT。
关键说明:非 nullafter 必须是完整对象,Connector 会用它完整替换同一 _key 的文档;after=null 会进入按 Key 删除的路径。CDC 模式会丢弃顶层 Value 为 null 的墓碑,因此删除事件必须由非 null 信封中的 after=null 表达。对于无 Schema 的对象,缺少 after 与显式 after=null 无法区分,也会被当作删除;在启用前应校验上游事件结构,并确保 Key 是受支持的单字段对象。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。由于 Connector 没有插件专属健康检查或指标,还应结合 ArangoDB 中的目标文档状态和 Worker 日志判断写入是否符合预期。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • Connector 不会自动创建 ArangoDB Database 或 Collection,也不提供显式的 Topic 到 Collection 映射;每条记录只能写入当前 Topic 最后一个句点后缀对应的预创建 Collection。
  • null Value 会作为完整文档插入或替换,不支持字段级局部更新、创建专用模式、软删除、外键转换或边关系映射;新 Value 中省略的旧字段不会被 Connector 合并保留。
  • Kafka Key 必须是至少包含一个字段且首字段值非 null 的对象,Connector 只使用第一个字段值;原始类型 Key、空 Key 和稳定的复合键映射不受支持。非 null Value 必须转换为对象,数组和标量不受支持。
  • 1.0.4 不提供 Connector 级 TLS、Truststore、连接超时、重试次数或退避配置。
  • 记录转换异常和 ArangoDB 写入异常会从 task.put 直接抛出并使 Task 失败;errors.tolerance 和 DLQ 不能处理此路径中的非法 Key、非法 Value 或数据库写入错误。
  • Connector 与 Kafka Offset 之间没有跨系统事务。一个轮询批次中的早期写入可能已经成功,而后续写入使 Task 失败;人工重启后,最后一次已提交 Offset 之后的记录可能重放。按稳定 _key 完整替换或删除可以降低重复创建文档的概率,但不构成 exactly-once、无重放或全局顺序保证。
  • 1.0.4 可能在 INFO 级别记录包含 arangodb.password 的原始配置 Map;仅限制日志读取并不能为传输链路提供加密,生产使用前需要单独评估凭证暴露风险。

常见问题

为什么更新后原文档中的某些字段消失了?

Connector 对非 null Value 执行完整文档替换,不做字段合并。检查新消息 Value 是否包含最终文档需要保留的所有字段,并确认 Key 的第一个字段仍指向同一 _key;如果上游只产生字段差异,应先在 Kafka Streams 或其他处理层重建完整对象,再交给 Connector 写入。

为什么记录写入了错误的 Collection,或因 Collection 不存在而失败?

Connector 使用 Topic 最后一个句点后的字符串作为 Collection 名称。检查实际 Topic 名称及 topicstopics.regex 的匹配范围,并预先创建对应 Collection;例如 dbserver1.inventory.customers 需要 customers Collection。Connector 不会通过配置把该 Topic 重映射到其他 Collection。

为什么墓碑记录没有删除文档?

先确认 Key 是至少包含一个字段且首字段值非 null 的对象。普通模式下,Value 为 null 的墓碑会按该 Key 删除;启用 CDC SMT 后,顶层墓碑会被丢弃,删除必须使用非 null CDC 信封并令 after=null。不要用缺少 after 的无 Schema 信封表示其他事件,因为它也会进入删除路径。

Task 失败后应如何恢复,是否会重复写入?

先从 Worker 日志和 ArangoDB 状态确认根因,修复非法记录、权限、目标资源、网络或服务端问题,再人工重启失败的 Task。普通 Connector 或驱动异常不会由插件自动分类并重试,也不能依靠 errors.tolerance 或 DLQ 跳过 task.put 内的失败。由于成功写入与 Kafka Offset 提交不在同一事务中,重启后可能重放最后一次已提交 Offset 之后的记录,包括失败前已经写入的早期批次;恢复后应核对目标文档和消费组 Offset,不要把 Task 恢复为 RUNNING 视为没有重放或数据覆盖。