概述
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。people Topic 的消息值能够提供 id、name 和 surname 字段;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 只能是
neo4j、neo4j+s、neo4j+ssc、bolt、bolt+s或bolt+ssc。多个 URI 会用于自定义地址解析;省略端口时使用7687。 - 必填:是
neo4j.database
写入的 Neo4j 数据库名称。
- 类型:
string - 默认值:
"" - 重要级别:高
- 有效值 / 注意事项:空字符串表示由 Neo4j 驱动或服务器选择数据库;建议显式指定目标数据库。
身份认证
neo4j.authentication.type
连接 Neo4j 时使用的认证类型。
- 类型:
string - 默认值:
BASIC - 重要级别:高
- 有效值 / 注意事项:可选
NONE、BASIC、KERBEROS、BEARER或CUSTOM;所选类型决定必须提供的凭证配置。
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
是否为普通 bolt 或 neo4j URI 启用由插件管理的加密。
- 类型:
string(逻辑类型为布尔值) - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:只能使用
true或false。+s和+sscURI 由 URI 自身控制加密,不使用此开关。
neo4j.security.trust-strategy
插件管理 TLS 时使用的证书信任策略。
- 类型:
string - 默认值:
TRUST_SYSTEM_CA_SIGNED_CERTIFICATES - 重要级别:低
- 有效值 / 注意事项:可选
TRUST_ALL_CERTIFICATES、TRUST_CUSTOM_CA_SIGNED_CERTIFICATES或TRUST_SYSTEM_CA_SIGNED_CERTIFICATES;仅在普通 URI 启用插件管理加密时生效。 - 必填:
conditional(启用插件管理加密时必须为非空值)
neo4j.security.hostname-verification-enabled
TLS 握手时是否校验证书主机名。
- 类型:
string(逻辑类型为布尔值) - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:只能使用
true或false;仅适用于普通 URI 上启用的插件管理加密。 - 必填:
conditional(启用插件管理加密时必须为非空值)
neo4j.security.cert-files
自定义 CA 证书文件列表。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:每个条目必须是 Worker 上可读的绝对路径普通文件;选择
TRUST_CUSTOM_CA_SIGNED_CERTIFICATES时列表不能为空。 - 必填:
conditional(使用自定义 CA 信任策略时必填)
连接池与事务重试
neo4j.connection-timeout
建立 Neo4j TCP 连接的超时时间。
- 类型:
string(逻辑类型为时长) - 默认值:
30s - 重要级别:低
- 有效值 / 注意事项:必须是非负组合时长,可使用
ms、s、m、h或d。
neo4j.pool.max-connection-pool-size
Neo4j 驱动连接池保留的最大连接数。
- 类型:
int - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:必须至少为
1;应结合 Task 并行度和目标数据库容量调整。
neo4j.pool.connection-acquisition-timeout
从连接池获取连接的最长等待时间。
- 类型:
string(逻辑类型为时长) - 默认值:
1m - 重要级别:低
- 有效值 / 注意事项:必须是非负组合时长,可使用
ms、s、m、h或d。
neo4j.pool.max-connection-lifetime
池中连接的最长生命周期。
- 类型:
string(逻辑类型为时长) - 默认值:
1h - 重要级别:低
- 有效值 / 注意事项:必须是非负组合时长,可使用
ms、s、m、h或d。
neo4j.pool.idle-time-before-connection-test
空闲连接在执行存活检查前可保持空闲的时长。
- 类型:
string(逻辑类型为时长) - 默认值:
"" - 重要级别:低
- 有效值 / 注意事项:可留空或使用非负组合时长;留空会保留驱动禁用或默认的存活检查语义。
neo4j.max-retry-time
Neo4j 驱动对可重试事务进行重试的最长时间。
- 类型:
string(逻辑类型为时长) - 默认值:
30s - 重要级别:低
- 有效值 / 注意事项:必须是非负组合时长,可使用
ms、s、m、h或d;它不控制 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 - 重要级别:低
- 有效值 / 注意事项:只能使用
true或false。设为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 - 重要级别:低
- 有效值 / 注意事项:只能使用
true或false;仅用于节点和关系 Pattern 策略。启用后,属性变化可能影响节点匹配结果。
neo4j.pattern.merge-relationship-properties
关系 Pattern 策略是否把传入属性加入关系 MERGE 匹配条件。
- 类型:
string(逻辑类型为布尔值) - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:只能使用
true或false;仅用于关系 Pattern 策略。启用后,属性变化可能影响关系匹配结果。
批处理与目标端 Offset
neo4j.batch-size
每个 Topic 生成一批写入语句时最多包含的事件数。
- 类型:
int - 默认值:
1000 - 重要级别:中
- 有效值 / 注意事项:必须至少为
1。存在 APOC 时它用于切分写事务批次;原生路径中它限制单条生成语句的事件数,但同一 Topic 组的多条语句仍可位于一个 Neo4j 托管事务中。
neo4j.batch-timeout
公开注册的批处理时长配置。
- 类型:
string(逻辑类型为时长) - 默认值:
0s - 重要级别:中
- 有效值 / 注意事项:接受使用
ms、s、m、h或d的非负组合时长;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 - 重要级别:中
- 有效值 / 注意事项:可选
none或all。all允许 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 节点的匹配键,并写入 name 和 surname。
id、name 和 surname 的 Map,并在目标数据库为业务唯一键建立合适的约束。Pattern 中的 !id 把 id 标记为节点匹配键,name 和 surname 作为写入属性;缺少有效唯一性时,可能匹配或创建非预期实体。不要让同一个 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。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.topics、neo4j.cdc.schema.topics、neo4j.cdc.source-id.topics 之一;删除重复分配和不在 topics 中的额外策略 Topic 后重启 Connector。
Connector 无法连接或认证 Neo4j
URI Scheme、数据库名、认证类型与凭证不匹配都可能导致连接失败。先确认neo4j.uri 可从 Worker 访问,再核对所选认证类型要求的非空凭证;使用 TLS 时同时检查 URI Scheme、信任策略、主机名校验和自定义 CA 文件的绝对路径及可读权限。
重启后出现重复节点或重复关系
Kafka Offset 提交可能晚于已经提交的 Neo4j 事务,非幂等 Cypher 或 CUDCREATE 因此可能在重放时产生重复结果。优先使用稳定业务键、约束和 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 失败。