Skip to main content

概述

Neo4j Source Connector 将 Neo4j 图数据发送到 Kafka,供下游搜索、分析和事件处理应用消费。QUERY 模式周期性执行 Cypher,将每一行查询结果写入一个 Topic,并用结果中的整数游标记录读取位置。CDC 模式读取数据库的变化日志,将匹配的节点和关系变化路由到配置的 Topic,可传递创建、更新和删除事件。 QUERY 适合按业务查询导出数据,再持续读取游标之后的记录;CDC 适合持续接收图变化。两种模式的最早起始位置含义不同:QUERY 使用数值游标起点,CDC 使用仍保留的变化历史起点。即使选择最早位置,CDC 也不会全量导出已有图数据。

前置条件

  • QUERY 使用的 Neo4j 账号须能够读取查询涉及的数据。
  • QUERY 查询须返回非空的 Long 类型游标,并以唯一、单调递增的游标设计支持分页和后续读取。
  • CDC 须使用支持 CDC 的 Neo4j 部署,在目标数据库启用变化捕获,并为账号授予读取变化日志的权限。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、Neo4j 数据库和目标 Topic,确认网络连通及访问权限;创建和管理操作参见管理 Connector。以下配置从数值起点读取已有记录,并继续轮询新记录。
替换连接地址、账号凭据和 Topic。源端 Event 节点须包含业务标识 id、Long 类型的 sequence 和待发送的 payloadsequence 须非负、唯一,后续写入不得使用已经读取过的序号。该字段返回为默认游标名 timestamp,但此处表示序号而非时间。账号默认使用 BASIC,数据库默认由驱动选择;若需指定数据库,设置 neo4j.database。默认一个 Task,Key 和 Value 的序列化沿用 Worker 的 Converter。 首次无可用 Offset 时,EARLIEST$lastCheck 设为 -1;有可用 Offset 时优先恢复。查询不提供一致性图快照,也不能自动观测硬删除。应用须在每次需要传递的变化发生时生成新序号;需要保留每次变化时,应写入独立的事件记录,而不是反复覆盖同一个节点。

配置

连接与数据库

neo4j.uri

配置 Neo4j 连接地址。
  • 类型LIST
  • 默认值:无,必填
  • 重要级别:高
  • 有效值 / 注意事项:非空 URI 列表,多个地址以逗号分隔。支持 neo4jneo4j+sneo4j+sscboltbolt+sbolt+ssc。首个 URI 用于创建驱动;配置多个 URI 时,全部地址均加入驱动的地址解析器。

neo4j.database

选择读取的数据库。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:留空或空白使用驱动的默认数据库选择;修改数据库会改变源 Offset 身份。

认证

neo4j.authentication.type

选择数据库认证方式。
  • 类型STRING
  • 默认值BASIC
  • 重要级别:高
  • 有效值 / 注意事项NONEBASICKERBEROSBEARERCUSTOM,区分大小写。对应方式的必需凭据不能为空。

neo4j.authentication.basic.username

配置 BASIC 账号。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:BASIC 时必填、非空。

neo4j.authentication.basic.password

配置 BASIC 密码。
  • 类型PASSWORD
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:BASIC 时必填、非空;不要将密码写入日志或共享配置文件。

neo4j.authentication.basic.realm

配置 BASIC 认证域。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:可选,按数据库认证设置填写。

neo4j.authentication.kerberos.ticket

配置 Kerberos 票据。
  • 类型PASSWORD
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:KERBEROS 时必填、非空;按敏感凭据管理。

neo4j.authentication.bearer.token

配置 Bearer Token。
  • 类型PASSWORD
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:BEARER 时必填、非空;按敏感凭据管理。

neo4j.authentication.custom.scheme

配置自定义认证方案。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:CUSTOM 时必填、非空,与服务端认证方案对应。

neo4j.authentication.custom.principal

配置自定义认证主体。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:CUSTOM 时必填、非空。

neo4j.authentication.custom.credentials

配置自定义认证凭据。
  • 类型PASSWORD
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:CUSTOM 时必填、非空;按敏感凭据管理。

neo4j.authentication.custom.realm

配置自定义认证域。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:CUSTOM 的可选项,与服务端配置对应。

加密与证书

neo4j.security.encrypted

控制普通 URI 连接的加密。
  • 类型STRING
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅接受 truefalse。URI 使用 +s+ssc 时由 URI 控制加密行为,此设置不覆盖它们。

neo4j.security.trust-strategy

选择显式加密连接的证书信任方式。
  • 类型STRING
  • 默认值TRUST_SYSTEM_CA_SIGNED_CERTIFICATES
  • 重要级别:低
  • 有效值 / 注意事项TRUST_ALL_CERTIFICATESTRUST_SYSTEM_CA_SIGNED_CERTIFICATESTRUST_CUSTOM_CA_SIGNED_CERTIFICATES。用于普通 URI 的显式加密,+s+ssc 使用自身信任行为;信任所有证书会削弱身份验证。

neo4j.security.hostname-verification-enabled

控制显式证书信任策略的主机名验证。
  • 类型STRING
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:仅接受 truefalse;使用显式信任策略时生效,不覆盖 URI 自带策略。

neo4j.security.cert-files

指定自定义 CA 证书文件。
  • 类型LIST
  • 默认值:空列表
  • 重要级别:低
  • 有效值 / 注意事项:使用自定义 CA 信任策略时提供 Worker 内可读的绝对文件路径,多个文件以逗号分隔;不能仅依赖界面提示判断是否需要证书。

连接池与重试

neo4j.connection-timeout

设置建立连接的超时时间。
  • 类型STRING
  • 默认值30s
  • 重要级别:低
  • 有效值 / 注意事项:非负整数配合小写单位 mssmhd,可组合如 1m30s;时间配置均应使用这种完整格式。

neo4j.pool.max-connection-pool-size

限制驱动连接池大小。
  • 类型INT
  • 默认值100
  • 重要级别:低
  • 有效值 / 注意事项:至少为 1;不是读取 Task 数量或安全分片数量。

neo4j.pool.connection-acquisition-timeout

设置从连接池获取连接的超时时间。
  • 类型STRING
  • 默认值1m
  • 重要级别:低
  • 有效值 / 注意事项:使用非负时长及小写单位,如 30s

neo4j.pool.max-connection-lifetime

设置连接的最大存活时间。
  • 类型STRING
  • 默认值1h
  • 重要级别:低
  • 有效值 / 注意事项:使用非负时长及小写单位。

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

设置连接空闲多久后在复用前检查连接。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:低
  • 有效值 / 注意事项:空值禁用该检查;非空时使用非负时长及小写单位。默认值不是数值 -1

neo4j.max-retry-time

设置驱动托管事务的最大重试时长。
  • 类型STRING
  • 默认值30s
  • 重要级别:低
  • 有效值 / 注意事项:使用非负时长及小写单位;不表示 Task 无限自动恢复,也不能代替故障排查和重启。

读取策略与起始位置

neo4j.source-strategy

选择读取模式。
  • 类型STRING
  • 默认值QUERY
  • 重要级别:高
  • 有效值 / 注意事项QUERYCDC。QUERY 需要查询和 Topic;CDC 至少需要一个 Topic Pattern,且不接受 RAW_JSON_STRING

neo4j.start-from

选择没有可复用 Offset 时的起始位置。
  • 类型STRING
  • 默认值NOW
  • 重要级别:高
  • 有效值 / 注意事项EARLIESTNOWUSER_PROVIDED。QUERY 的 EARLIEST 为 -1,NOW 为本机当前 epoch 毫秒;序号游标不应使用 NOW。CDC 分别使用最早保留 ID、当前 ID 或提供的 ID。有可用 Offset 时优先恢复。

neo4j.start-from.value

提供自定义起始游标。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:USER_PROVIDED 时必填。QUERY 为可解析的有符号 Long 字符串,CDC 为该数据库仍有效的变化 ID。使用严格大于条件的查询不包含起始游标对应的记录。

neo4j.ignore-stored-offset

控制是否忽略已保存 Offset。
  • 类型STRING
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:仅接受 truefalse。true 强制按 start-from 选择位置,可能重放或跳过数据;每次重启都会再次应用,不宜作为日常恢复设置。

查询与游标

neo4j.query

配置需要执行的 Cypher 查询。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:QUERY 时必填。通过 $lastCheck 接收当前游标;查询须自行过滤、升序排列并返回 Long 游标,Connector 不自动添加 WHERE 或 ORDER BY。修改查询文本会改变源 Offset 身份。

neo4j.query.streaming-property

指定结果行中用于推进游标的字段名。
  • 类型STRING
  • 默认值timestamp
  • 重要级别:低
  • 有效值 / 注意事项:非空白字段名,结果中该字段须为 Long。修改字段会使已有 Offset 不可复用。游标须唯一,结果须按游标升序返回;仅排序不能避免严格大于分页在相同游标处遗漏记录。

neo4j.query.topic

指定查询结果的单个目标 Topic。
  • 类型STRING
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:QUERY 时必填;配置校验不检查 Topic 是否存在。更换 Topic 不会自动重置源 Offset。

neo4j.query.timeout

设置数据库事务超时时间。
  • 类型STRING
  • 默认值0s
  • 重要级别:低
  • 有效值 / 注意事项:正时长覆盖事务超时,0 不覆盖数据库设置;共享事务配置也可影响 CDC 事务。不是 poll-duration 的硬截止时间。

neo4j.query.poll-interval

设置 QUERY 无数据时再次查询的等待间隔。
  • 类型STRING
  • 默认值1s
  • 重要级别:中
  • 有效值 / 注意事项:非负时长;更短间隔增加空查询压力,不改变 Offset 提交频率。

neo4j.query.poll-duration

设置一次 QUERY poll 的循环时间预算。
  • 类型STRING
  • 默认值5s
  • 重要级别:中
  • 有效值 / 注意事项:应使用正时长,0 不提供读取时间预算。仅在操作之间检查,不会强制中断阻塞查询或等待。

neo4j.query.force-maps-as-struct

控制 QUERY 结果中的 Map 是否转为 Struct。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项truefalse;仅适用于 QUERY,须与消费端预期的数据结构一致。

批量与负载格式

neo4j.batch-size

限制每批读取的数量。
  • 类型INT
  • 默认值1000
  • 重要级别:中
  • 有效值 / 注意事项:至少为 1。QUERY 限制结果行,CDC 限制变化事件;CDC 路由展开后的消息数可能更大。不是字节限制,增大时需关注内存和下游压力。

neo4j.payload-mode

选择图数据的负载表示方式。
  • 类型STRING
  • 默认值EXTENDED
  • 重要级别:中
  • 有效值 / 注意事项EXTENDEDCOMPACTRAW_JSON_STRING。EXTENDED 使用类型包装,COMPACT 更紧凑但需关注类型变化;RAW_JSON_STRING 仅用于 QUERY。负载模式不替代 Converter,也不自动提供 Schema Registry 或 Schema 演进保障。

CDC 轮询与读取选项

neo4j.cdc.use-leader

配置 CDC 的 Leader 读取选项。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse。不要将此选项作为强制 CDC 路由到 Leader 的保障;如需固定读取节点,应结合部署路由确认实际连接行为。

neo4j.cdc.poll-interval

设置 CDC 无数据时的等待间隔。
  • 类型STRING
  • 默认值1s
  • 重要级别:中
  • 有效值 / 注意事项:非负时长,按延迟需求和数据库查询负载取舍。

neo4j.cdc.poll-duration

设置一次 CDC poll 的循环时间预算。
  • 类型STRING
  • 默认值5s
  • 重要级别:中
  • 有效值 / 注意事项:应使用正时长;预算仅在操作之间检查,不是数据库调用的硬超时。

CDC Topic 路由与事件表示

neo4j.cdc.topic.<topic>.patterns

为目标 Topic 设置节点或关系选择 Pattern。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:CDC 动态配置,至少提供一个实际 Pattern;例如 (:Event) 匹配 Event 节点。Topic 名支持字母、数字、点、下划线和连字符。同一 Topic 不可与索引形式混用;Pattern 投影会改变输出属性。

neo4j.cdc.topic.<topic>.patterns.<index>.pattern

按索引配置单个 CDC 选择 Pattern。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:索引从 0 开始且连续,每个索引一个 Pattern;同一 Topic 不可与非索引形式混用。

neo4j.cdc.topic.<topic>.patterns.<index>.operation

按变化操作过滤该索引的事件。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项createupdatedelete,运行时转为小写;须存在相同索引的 Pattern。省略时不设置操作过滤。

neo4j.cdc.topic.<topic>.patterns.<index>.changesTo

按发生变化的属性过滤事件。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:逗号分隔的属性名,会去除两端空白;须存在相同索引的 Pattern。

neo4j.cdc.topic.<topic>.patterns.<index>.metadata.<metadata-key>

按事件元数据过滤该索引的变化。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值
  • 重要级别:未声明
  • 有效值 / 注意事项:metadata-key 支持 authenticatedUserexecutingUsertxMetadata.<key>;键名支持字母、数字、点、下划线和连字符。须存在相同索引的 Pattern。

neo4j.cdc.topic.<topic>.key-strategy

选择 CDC 消息 Key 的构造方式。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值;未设置时运行时使用 WHOLE_VALUE
  • 重要级别:未声明
  • 有效值 / 注意事项SKIPELEMENT_IDENTITY_KEYSWHOLE_VALUE,区分大小写。ENTITY_KEYS 在实体没有键时输出 null;不要据此假定所有事件都有业务键。

neo4j.cdc.topic.<topic>.value-strategy

选择 CDC 消息 Value 的构造方式。
  • 类型STRING
  • 默认值:无 ConfigDef 默认值;未设置时运行时使用 CHANGE_EVENT
  • 重要级别:未声明
  • 有效值 / 注意事项CHANGE_EVENTENTITY_EVENT,区分大小写;前者保留变化事件表示,后者使用实体事件表示。

CDC 指标采集选项

neo4j.cdc.metric.last-db-tx-id.enabled

控制是否采集数据库最新事务 ID 指标。
  • 类型STRING
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅接受 truefalse;启用后会增加相应刷新查询。

neo4j.cdc.metric.last-db-tx-id.refresh-interval

设置最新事务 ID 指标刷新间隔。
  • 类型STRING
  • 默认值30s
  • 重要级别:低
  • 有效值 / 注意事项:非负时长,仅在该指标启用时有意义。

Connector 运行与序列化

connector.class

指定 Source Connector 实现类。
  • 类型STRING
  • 默认值:无,必填
  • 重要级别:高
  • 有效值 / 注意事项:使用 org.neo4j.connectors.kafka.source.Neo4jConnector

tasks.max

设置 Task 数量上限。
  • 类型INT
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:至少为 1;单条配置流保持一个 Task。多个 Task 会独立读取相同源数据且共享逻辑 Offset 身份,不实现安全的源端分片。

key.converter

指定消息 Key 的序列化 Converter。
  • 类型CLASS
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:留空继承 Worker 配置;设置时使用可实例化的 Converter 类,如 org.apache.kafka.connect.json.JsonConverter,并按其要求配置序列化选项。

value.converter

指定消息 Value 的序列化 Converter。
  • 类型CLASS
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:留空继承 Worker 配置;设置时使用可实例化的 Converter 类。须能处理所选 payload-mode 产生的数据类型。

最佳实践

持续接收图数据的创建、更新和删除

适用业务场景:已有数据库需要把后续图变化发送给下游应用,尤其需要传递硬删除时,使用 CDC 而不是查询当前节点状态。开始前先启用 CDC;已有实体的基线应另外规划导出,不能依赖 CDC 自动补齐。 配置示例:沿用快速开始的连接与认证配置,移除 neo4j.queryneo4j.query.topic,覆盖 neo4j.start-from 并添加以下配置。将动态键中的 <cdc-topic> 替换为真实 Topic 名。
关键说明:在无可复用 Offset 的新 Connector 上,NOW 从当前 CDC 位置读取之后的 Event 节点变化;运行后继续使用保存的 Offset。基线与变化流的衔接应在源端保留期内设计,避免导出期间的变化缺口。该 Pattern 不限制操作类型,匹配的创建、更新和删除均可发送;默认使用 CHANGE_EVENT。此配置不自动启用 exactly-once,普通运行下消费端应能处理故障恢复后的重放。

维护后从指定位置重放查询记录

适用业务场景:使用 QUERY 的流需要补发某个已知游标之后的数据,例如修复下游处理问题后重新处理记录;源端仍保留所需事件,消费端已准备好去重或幂等处理。 配置示例:在快速开始的 QUERY 配置上覆盖起始位置并添加以下配置。<last-processed-sequence> 为希望重放范围之前的最后一个序号;这里是源端数值位置,不是 Kafka Offset。
关键说明:暂停 Connector 后再修改配置,查询的严格大于条件只读取指定序号之后的记录。重放可能产生重复消息,且只能读取当前仍存在的记录,不能恢复已硬删除的数据。确认消息发送且新 Offset 已持久化后,移除 neo4j.ignore-stored-offsetneo4j.start-from.value,将 neo4j.start-from 恢复为 EARLIEST,再恢复日常运行。恢复前保持数据库、查询文本和游标字段不变,避免重新选择起点。

监控

监控内容

关注 Kafka Connect 集群健康、Connector 和 Task 状态、吞吐、处理延迟、Offset 提交、错误及重试,同时观察 Worker JVM 的堆内存、GC 和线程状态。状态正常但无消息时,应结合源端数据变化和消费结果判断是否确实空闲;只有部署启用了相应错误处理时才关注 DLQ 活动。

导入 Grafana 大盘

下载 AutoMQ Connect Cluster Dashboard,确认 Kafka Connect 指标已接入 Prometheus 兼容数据源且集群、Worker、Connector 等标签与大盘查询一致,然后在 Grafana 中导入 JSON 并选择相应数据源。

限制条件

  • QUERY 只保存单个 Long 游标,不提供相同游标的二级分页位置;严格大于条件在批次边界存在相同游标时可能遗漏记录。
  • QUERY 不自动生成过滤或排序条件,也不能自动捕获硬删除。
  • CDC 的 EARLIEST 只读取仍保留的变化历史,不导出所有已有实体。
  • CDC 变化 ID 受数据库保留期和数据库历史影响,备份恢复、快照恢复以及 Aura 暂停后恢复可能使已有 ID 失效,不能依赖自动回退。
  • CDC 同一 Topic 不可混用列表和索引 Pattern,索引须从零连续编号。
  • 重叠的 CDC 选择器可为同一变化产生多条消息,不执行跨选择器去重。
  • 多 Task 不提供源端分片;增加 Task 数量可能导致重复读取及 Offset 干扰。
  • QUERY 不支持 Source exactly-once;CDC 的支持声明仍要求分布式 Worker 启用该能力、Broker 支持事务并正确授权,消费端使用事务隔离读取,不能视为自动具备端到端保障。
  • 不保证跨 Kafka 分区、Topic 或多个 Task 的全局顺序;CDC 同事务事件的序号也不保证对应原始操作执行顺序。
  • CDC 不接受 RAW_JSON_STRING 负载模式。

常见问题

Task 正常运行但没有新消息怎么办?

QUERY 先检查目标数据库、Topic 和查询结果,确认返回字段包含配置的 Long 游标,并且新记录的游标大于保存位置。默认 NOW 使用当前 epoch 毫秒,若业务字段是普通序号,应在首次启动时改为 EARLIEST;已有 Offset 时只修改 start-from 不会重置位置。CDC 检查数据库是否启用 CDC、账号权限和实际 Pattern 是否匹配,NOW 不会补发启动前的已有实体。需要重放时先确认源数据保留情况,再显式选择起点。

查询读取后有记录遗漏怎么办?

检查查询是否按游标升序返回、游标是否唯一,以及后写入记录是否使用了已越过的游标。Connector 按该批最后一行推进游标,时间戳相同或无序结果可能产生缺口。改用唯一单调序号或独立事件记录;已遗漏但仍保留的记录可在下游具备去重能力后从较早位置重放。仅添加 ORDER BY 无法解决相同游标的分页边界。

重启后出现重复消息怎么办?

普通 Source 运行中,消息发布成功但 Offset 尚未持久化时发生故障,恢复可能再次读取。检查是否长期设置 ignore-stored-offset=true、是否改变了数据库、查询文本或游标字段,以及是否部署了多个相同读取任务。保持一个 Task、稳定源身份并恢复默认 Offset 复用;消费端使用业务标识及事件版本去重。CDC 的 exactly-once 需由平台确认 Worker、Broker 和消费端的事务条件,不能仅设置 source-strategy 就获得保障。

CDC 恢复时报变化 ID 无效怎么办?

检查停机时长与 CDC 保留期,以及数据库是否发生恢复或 Aura 暂停后恢复。失效位置不能再用于连续读取;先明确缺失范围和重新建立基线的方案,再选择有效的 USER_PROVIDED ID、最早保留位置或当前位置。强制忽略旧 Offset 会改变读取范围,不能修复已经过期的历史。

更换 Topic 或扩大 CDC 选择范围后没有补齐旧数据怎么办?

Topic、连接 URI 和 CDC 选择器不参与相应源 Offset 身份。原 Connector 可能继续使用先前位置,不会因为输出范围改变而自动回填。先确定要保留的读取位置和需要补齐的范围;若需要从历史重放,确认保留期、去重方案和基线衔接后再显式调整起点,不要直接假定换 Topic 就会重新导出。

Task 因查询字段类型或认证错误失败怎么办?

检查完整异常日志但不要输出凭据。QUERY 的游标必须存在且为 Long,修正 Cypher 返回字段及源端类型后再恢复;认证失败时核对所选认证类型及其必需凭据,证书失败时检查 URI 加密方式和 Worker 内证书路径。最大事务重试时长不是所有 Task 错误的自动恢复策略,修正配置后通过 Connector 管理操作重启失败任务。