Skip to main content

概述

Neo4j Sink Connector 消费 Kafka Topic 中的记录,并把记录转换为对 Neo4j 或 Aura 数据库的节点、关系和属性写入。Connector 位于 Kafka 与图数据库之间,可按 Topic 选择 Cypher、Pattern、CUD、CDC Schema 或 CDC Source Id 策略;每个 Topic 必须且只能对应一种策略。 Cypher 策略适合用自定义语句控制图写入,Pattern 策略适合声明字段到节点或关系的映射,CUD 和 CDC 策略则用于消费相应的操作或变更事件格式。典型用途包括把业务事件构造成知识图谱、持续更新实体关系,以及在 Neo4j 数据库之间传递 CDC 事件。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和 Neo4j 目标数据库,并确认网络连通和访问权限。具体准备和管理操作请参阅管理 Connector
替换 Neo4j 地址、数据库和凭证占位符。示例要求 people Topic 的消息值能够提供 idnamesurname 字段;Connector 使用默认的 __value 绑定执行 Cypher,并按 id 合并 Person 节点。生产环境应通过安全的凭证管理方式提供密码。

配置

Connector 身份与 Topic 订阅

connector.class

要加载的 Neo4j Sink Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 org.neo4j.connectors.kafka.sink.Neo4jConnector
  • 必填:是

topics

要消费的 Kafka Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔的字面 Topic 名称。每个名称必须由且仅由一种 Sink 策略接管;Neo4j 5.5.1 不能使用 topics.regex 建立所需的逐 Topic 策略映射。
  • 必填:是

Neo4j 连接

neo4j.uri

一个或多个 Neo4j 连接 URI。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:URI Scheme 只能是 neo4jneo4j+sneo4j+sscboltbolt+sbolt+ssc。多个 URI 会用于自定义地址解析;省略端口时使用 7687
  • 必填:是

neo4j.database

写入的 Neo4j 数据库名称。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项:空字符串表示由 Neo4j 驱动或服务器选择数据库;建议显式指定目标数据库。

身份认证

neo4j.authentication.type

连接 Neo4j 时使用的认证类型。
  • 类型string
  • 默认值BASIC
  • 重要级别:高
  • 有效值 / 注意事项:可选 NONEBASICKERBEROSBEARERCUSTOM;所选类型决定必须提供的凭证配置。

neo4j.authentication.basic.username

Basic 认证用户名。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=BASIC 时必须为非空值。
  • 必填conditional(使用 Basic 认证时必填)

neo4j.authentication.basic.password

Basic 认证密码。
  • 类型password
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=BASIC 时必须为非空值;建议通过 Config Provider 或等效的密钥管理方式提供。
  • 必填conditional(使用 Basic 认证时必填)

neo4j.authentication.basic.realm

Basic 认证 Realm。
  • 类型string
  • 默认值""
  • 重要级别:低
  • 有效值 / 注意事项:仅用于 Basic 认证;空字符串选择默认 Realm。

neo4j.authentication.kerberos.ticket

Kerberos 认证票据。
  • 类型password
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=KERBEROS 时必须为非空值,并应通过安全的凭证管理方式提供。
  • 必填conditional(使用 Kerberos 认证时必填)

neo4j.authentication.bearer.token

Bearer 认证令牌。
  • 类型password
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=BEARER 时必须为非空值,并应通过安全的凭证管理方式提供。
  • 必填conditional(使用 Bearer 认证时必填)

neo4j.authentication.custom.scheme

自定义认证 Scheme。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=CUSTOM 时必须为非空值。
  • 必填conditional(使用自定义认证时必填)

neo4j.authentication.custom.principal

自定义认证 Principal。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=CUSTOM 时必须为非空值。
  • 必填conditional(使用自定义认证时必填)

neo4j.authentication.custom.credentials

自定义认证凭证。
  • 类型password
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项neo4j.authentication.type=CUSTOM 时必须为非空值,并应通过安全的凭证管理方式提供。
  • 必填conditional(使用自定义认证时必填)

neo4j.authentication.custom.realm

传给自定义 Neo4j Auth Token 的 Realm。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项:仅用于自定义认证;是否填写取决于认证提供方。

TLS

neo4j.security.encrypted

是否为普通 boltneo4j URI 启用由插件管理的加密。
  • 类型string(逻辑类型为布尔值)
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只能使用 truefalse+s+ssc URI 由 URI 自身控制加密,不使用此开关。

neo4j.security.trust-strategy

插件管理 TLS 时使用的证书信任策略。
  • 类型string
  • 默认值TRUST_SYSTEM_CA_SIGNED_CERTIFICATES
  • 重要级别:低
  • 有效值 / 注意事项:可选 TRUST_ALL_CERTIFICATESTRUST_CUSTOM_CA_SIGNED_CERTIFICATESTRUST_SYSTEM_CA_SIGNED_CERTIFICATES;仅在普通 URI 启用插件管理加密时生效。
  • 必填conditional(启用插件管理加密时必须为非空值)

neo4j.security.hostname-verification-enabled

TLS 握手时是否校验证书主机名。
  • 类型string(逻辑类型为布尔值)
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:只能使用 truefalse;仅适用于普通 URI 上启用的插件管理加密。
  • 必填conditional(启用插件管理加密时必须为非空值)

neo4j.security.cert-files

自定义 CA 证书文件列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:低
  • 有效值 / 注意事项:每个条目必须是 Worker 上可读的绝对路径普通文件;选择 TRUST_CUSTOM_CA_SIGNED_CERTIFICATES 时列表不能为空。
  • 必填conditional(使用自定义 CA 信任策略时必填)

连接池与事务重试

neo4j.connection-timeout

建立 Neo4j TCP 连接的超时时间。
  • 类型string(逻辑类型为时长)
  • 默认值30s
  • 重要级别:低
  • 有效值 / 注意事项:必须是非负组合时长,可使用 mssmhd

neo4j.pool.max-connection-pool-size

Neo4j 驱动连接池保留的最大连接数。
  • 类型int
  • 默认值100
  • 重要级别:低
  • 有效值 / 注意事项:必须至少为 1;应结合 Task 并行度和目标数据库容量调整。

neo4j.pool.connection-acquisition-timeout

从连接池获取连接的最长等待时间。
  • 类型string(逻辑类型为时长)
  • 默认值1m
  • 重要级别:低
  • 有效值 / 注意事项:必须是非负组合时长,可使用 mssmhd

neo4j.pool.max-connection-lifetime

池中连接的最长生命周期。
  • 类型string(逻辑类型为时长)
  • 默认值1h
  • 重要级别:低
  • 有效值 / 注意事项:必须是非负组合时长,可使用 mssmhd

neo4j.pool.idle-time-before-connection-test

空闲连接在执行存活检查前可保持空闲的时长。
  • 类型string(逻辑类型为时长)
  • 默认值""
  • 重要级别:低
  • 有效值 / 注意事项:可留空或使用非负组合时长;留空会保留驱动禁用或默认的存活检查语义。

neo4j.max-retry-time

Neo4j 驱动对可重试事务进行重试的最长时间。
  • 类型string(逻辑类型为时长)
  • 默认值30s
  • 重要级别:低
  • 有效值 / 注意事项:必须是非负组合时长,可使用 mssmhd;它不控制 Kafka Connect 框架的错误容忍或 DLQ 行为。

写入策略与 Topic 路由

neo4j.cypher.topic.<topic>

为指定 Topic 执行的 Cypher 语句。
  • 类型string
  • 默认值:无
  • 重要级别:N/A
  • 有效值 / 注意事项:将 <topic> 替换为 topics 中的精确 Topic 名称;每个 Topic 只能配置一种策略。Cypher 在任务运行时解析和执行。
  • 必填conditional(该 Topic 使用 Cypher 策略时必填)

neo4j.pattern.topic.<topic>

指定 Topic 的节点或关系 Pattern 映射。
  • 类型string
  • 默认值:无
  • 重要级别:N/A
  • 有效值 / 注意事项:将 <topic> 替换为 topics 中的精确 Topic 名称;Pattern 必须能解析为节点或关系映射,且该 Topic 不能再选择其他策略。
  • 必填conditional(该 Topic 使用 Pattern 策略时必填)

neo4j.cud.topics

使用 CUD 操作格式的 Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:列表中的 Topic 必须同时出现在 topics,且不能分配给其他策略;消息值必须符合 Connector 的 CUD 结构。

neo4j.cdc.schema.topics

使用 CDC Schema 策略的 Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:列表中的 Topic 必须同时出现在 topics,且不能分配给其他策略;消息必须是受支持的 Neo4j CDC Schema 或兼容的旧 Streams 事件结构。

neo4j.cdc.source-id.topics

使用 CDC Source Id 策略的 Topic 列表。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:列表中的 Topic 必须同时出现在 topics,且不能分配给其他策略;事件必须包含源元素 ID。

neo4j.cdc.source-id.label-name

CDC Source Id 策略为目标节点使用的身份 Label。
  • 类型string
  • 默认值SourceEvent
  • 重要级别:低
  • 有效值 / 注意事项:仅在 neo4j.cdc.source-id.topics 非空时使用;应为该来源稳定保留,避免不同来源的元素 ID 发生碰撞。

neo4j.cdc.source-id.property-name

CDC Source Id 策略保存源元素 ID 的属性名。
  • 类型string
  • 默认值sourceId
  • 重要级别:低
  • 有效值 / 注意事项:仅在 neo4j.cdc.source-id.topics 非空时使用;更改属性名会改变目标身份匹配方式。

Cypher 记录绑定

neo4j.cypher.bind-timestamp-as

在 Cypher 中暴露 Kafka 记录时间戳的变量名。
  • 类型string
  • 默认值__timestamp
  • 重要级别:中
  • 有效值 / 注意事项:空字符串会禁用该绑定;时间戳转换为 UTC Offset Date-Time。禁用绑定时仍必须保留至少一个有效 Cypher 绑定。

neo4j.cypher.bind-header-as

在 Cypher 中暴露 Kafka Header 的变量名。
  • 类型string
  • 默认值__header
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该绑定;禁用绑定时仍必须保留至少一个有效 Cypher 绑定。

neo4j.cypher.bind-key-as

在 Cypher 中暴露 Kafka 记录 Key 的变量名。
  • 类型string
  • 默认值__key
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该绑定;禁用绑定时仍必须保留至少一个有效 Cypher 绑定。

neo4j.cypher.bind-value-as

在 Cypher 中暴露 Kafka 记录 Value 的变量名。
  • 类型string
  • 默认值__value
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该绑定;禁用绑定时仍必须保留至少一个有效 Cypher 绑定。

neo4j.cypher.bind-value-as-event

是否同时使用旧版 event 变量名绑定记录 Value。
  • 类型string(逻辑类型为布尔值)
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:只能使用 truefalse。设为 false 时,至少一个显式 Cypher 绑定必须保持非空。

Pattern 记录绑定与属性写入

neo4j.pattern.bind-timestamp-as

Pattern 表达式使用的 Kafka 记录时间戳别名。
  • 类型string
  • 默认值__timestamp
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该别名。

neo4j.pattern.bind-header-as

Pattern 表达式使用的 Kafka Header 别名。
  • 类型string
  • 默认值__header
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该别名。

neo4j.pattern.bind-key-as

Pattern 表达式使用的 Kafka 记录 Key 别名。
  • 类型string
  • 默认值__key
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该别名。

neo4j.pattern.bind-value-as

Pattern 表达式使用的 Kafka 记录 Value 别名。
  • 类型string
  • 默认值__value
  • 重要级别:低
  • 有效值 / 注意事项:空字符串会禁用该别名;非 Tombstone 消息的 Value 必须可转换为 Map。

neo4j.pattern.merge-node-properties

Pattern 策略是否把传入属性加入节点 MERGE 匹配条件。
  • 类型string(逻辑类型为布尔值)
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只能使用 truefalse;仅用于节点和关系 Pattern 策略。启用后,属性变化可能影响节点匹配结果。

neo4j.pattern.merge-relationship-properties

关系 Pattern 策略是否把传入属性加入关系 MERGE 匹配条件。
  • 类型string(逻辑类型为布尔值)
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只能使用 truefalse;仅用于关系 Pattern 策略。启用后,属性变化可能影响关系匹配结果。

批处理与目标端 Offset

neo4j.batch-size

每个 Topic 生成一批写入语句时最多包含的事件数。
  • 类型int
  • 默认值1000
  • 重要级别:中
  • 有效值 / 注意事项:必须至少为 1。存在 APOC 时它用于切分写事务批次;原生路径中它限制单条生成语句的事件数,但同一 Topic 组的多条语句仍可位于一个 Neo4j 托管事务中。

neo4j.batch-timeout

公开注册的批处理时长配置。
  • 类型string(逻辑类型为时长)
  • 默认值0s
  • 重要级别:中
  • 有效值 / 注意事项:接受使用 mssmhd 的非负组合时长;Neo4j Connector 5.5.1 的 Sink 主执行路径没有读取此配置,因此不能依靠它限制批次或事务执行时间。

neo4j.max-batched-queries

原生批处理路径中一批可组合的不同语句形状上限。
  • 类型int
  • 默认值50
  • 重要级别:低
  • 有效值 / 注意事项:必须至少为 1;达到上限时原生路径会切分生成的工作,APOC 路径不读取此配置。

neo4j.eos-offset-label

保存目标端 Offset 节点时使用的 Neo4j Label。
  • 类型string
  • 默认值""
  • 重要级别:高
  • 有效值 / 注意事项:空字符串禁用目标端 Offset 节点。非空时,Connector 按策略、Topic 和分区保存已写入 Offset,并要求目标数据库具备相应约束;该机制只在相同身份范围内过滤已提交到 Neo4j 的 Offset,不构成端到端或全局 exactly-once 保证。同一 Task 一次处理同 Topic 多分区记录时不应依赖此机制。

Task 与 Converter

tasks.max

Connector 请求创建的 Task 数量上限。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须至少为 1;实际并行度受 Kafka 分区分配和 Neo4j 并发容量限制,跨 Task、Topic 或分区没有全局顺序。

tasks.max.enforce

是否强制 Connector 返回的 Task 数不超过 tasks.max
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:Kafka Connect 3.9.1 已弃用此配置并计划在未来主要版本移除,未提供替代项;建议保持 true
  • 已弃用:是

key.converter

Connector 级 Kafka 记录 Key Converter 类。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 的 Key Converter;显式值必须是可实例化的 Kafka Connect Converter 实现。Converter 专属配置由相应插件定义。

value.converter

Connector 级 Kafka 记录 Value Converter 类。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 的 Value Converter;显式值必须是可实例化的 Kafka Connect Converter 实现。所有 Sink 策略都会使用转换后的 Value。

错误处理与 DLQ

errors.tolerance

Kafka Connect 对记录错误的容忍策略。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 noneallall 允许 Errant Record Reporter 隔离不可处理记录,但单独设置此项不会配置 DLQ。该机制不同于 neo4j.max-retry-time 控制的数据库事务重试。

errors.deadletterqueue.topic.name

写入不可处理记录的 DLQ Topic 名称。
  • 类型string
  • 默认值""
  • 重要级别:中
  • 有效值 / 注意事项:空字符串禁用 DLQ;非空 Topic 不能与业务订阅 Topic 相同。应与 errors.tolerance=all 配合使用。

errors.deadletterqueue.topic.replication.factor

Kafka Connect 创建 DLQ Topic 时使用的副本因子。
  • 类型short
  • 默认值3
  • 重要级别:中
  • 有效值 / 注意事项:必须适合 Kafka 集群的 Broker 数量;未配置 DLQ 或 Topic 已存在时不会用于创建 Topic。

errors.deadletterqueue.context.headers.enable

是否在 DLQ 记录中加入错误上下文 Header。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅在记录实际写入 DLQ 时生效;设为 true 可保留排障所需的 Connector 错误上下文。

最佳实践

首次接入时声明字段到图模型的映射

适用业务场景:首次把结构稳定的实体事件写入 Neo4j,希望不维护自定义 Cypher,并明确哪些字段用于识别节点、哪些字段作为节点属性。以下配置使用 Pattern 策略,以 id 作为 Person 节点的匹配键,并写入 namesurname
关键说明:先确保消息值可转换为包含 idnamesurname 的 Map,并在目标数据库为业务唯一键建立合适的约束。Pattern 中的 !idid 标记为节点匹配键,namesurname 作为写入属性;缺少有效唯一性时,可能匹配或创建非预期实体。不要让同一个 people Topic 同时配置 Cypher、CUD 或 CDC 策略。

持续写入时按分区扩展 Task

适用业务场景:Connector 已持续写入 Neo4j,单 Task 的吞吐或延迟不能满足要求,并且业务 Topic 有多个分区可供并行分配。以下配置在快速开始基础上把 Task 上限提高到 4
关键说明:tasks.max=4 只是请求的 Task 上限,不保证一定运行四个 Task;有效并行度不会超过可分配分区,也受 Neo4j 连接池和数据库写入容量限制。扩容前后应观察 Task 分配、吞吐、延迟和数据库负载。只依赖单个 Topic Partition 内的顺序,不要假定跨 Task 或跨分区存在全局顺序。

日常运行时用 DLQ 隔离坏记录

适用业务场景:持续写入期间可能出现少量无法转换、违反策略结构或导致 Neo4j 写入失败的记录,希望保留这些记录供排查,同时让可处理记录继续写入。以下配置启用错误容忍、DLQ 和上下文 Header。
关键说明:DLQ Topic 不能与 people 重叠,副本因子必须适合 Kafka 集群。Connector 会在批次失败后尝试隔离单条记录;成功记录可能已经提交,而被 Reporter 接受并写入 DLQ 的记录不会写入 Neo4j。此时 put 可以成功返回,后续提交的 Kafka Offset 可能越过这些记录。应持续监控 DLQ、修复数据或映射,并按业务要求单独回放;该配置不是事务重试或 exactly-once 保证。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

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

限制条件

  • Neo4j Connector 5.5.1 的 Sink 策略依赖 topics 中的精确 Topic 名称,不能使用 topics.regex 代替字面 Topic 列表。
  • 每个 Topic 必须且只能分配给 Cypher、Pattern、CUD、CDC Schema 或 CDC Source Id 中的一种策略;Connector 不支持按单条记录动态选择策略或回退策略。
  • neo4j.batch-timeout 虽可通过配置校验,但 5.5.1 的 Sink 主执行路径没有读取它,不能用于限制批次或事务执行时间。
  • Connector 不提供跨 Task、Topic 或 Kafka Partition 的全局写入顺序;CDC 事务 ID 和序号也不会重建源事务边界。
  • neo4j.eos-offset-label 只在稳定的策略、Topic、分区、目标数据库和 Label 身份内过滤已提交到 Neo4j 的 Offset,不构成端到端或全局 exactly-once 保证;同 Topic 多分区记录进入同一 Task 批次时不应依赖该机制。

常见问题

Task 启动时报 Topic 没有策略或分配了多个策略

topics 与策略配置不一致会阻止任务启动。检查每个字面 Topic 是否恰好出现在一个位置:对应的 neo4j.cypher.topic.<topic>neo4j.pattern.topic.<topic>,或 neo4j.cud.topicsneo4j.cdc.schema.topicsneo4j.cdc.source-id.topics 之一;删除重复分配和不在 topics 中的额外策略 Topic 后重启 Connector。

Connector 无法连接或认证 Neo4j

URI Scheme、数据库名、认证类型与凭证不匹配都可能导致连接失败。先确认 neo4j.uri 可从 Worker 访问,再核对所选认证类型要求的非空凭证;使用 TLS 时同时检查 URI Scheme、信任策略、主机名校验和自定义 CA 文件的绝对路径及可读权限。

重启后出现重复节点或重复关系

Kafka Offset 提交可能晚于已经提交的 Neo4j 事务,非幂等 Cypher 或 CUD CREATE 因此可能在重放时产生重复结果。优先使用稳定业务键、约束和 MERGE 设计幂等写入;若评估使用 neo4j.eos-offset-label,还必须保持策略、Topic、分区、目标数据库和 Label 不变。该机制不是全局 exactly-once,也不应在同一 Task 一次处理同 Topic 多分区记录时用作重放过滤保障。

坏记录没有进入 DLQ 或 Task 仍然失败

只设置 DLQ Topic 不会启用容错。确认同时设置了 errors.tolerance=all,DLQ Topic 与业务订阅不重叠,并且 Kafka Connect 有权访问或创建该 Topic;随后检查 DLQ 副本因子、Reporter 错误和 Neo4j 异常。不能被 Reporter 接受的错误仍会使 put 失败。