概述
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。以下配置从数值起点读取已有记录,并继续轮询新记录。Event 节点须包含业务标识 id、Long 类型的 sequence 和待发送的 payload;sequence 须非负、唯一,后续写入不得使用已经读取过的序号。该字段返回为默认游标名 timestamp,但此处表示序号而非时间。账号默认使用 BASIC,数据库默认由驱动选择;若需指定数据库,设置 neo4j.database。默认一个 Task,Key 和 Value 的序列化沿用 Worker 的 Converter。
首次无可用 Offset 时,EARLIEST 将 $lastCheck 设为 -1;有可用 Offset 时优先恢复。查询不提供一致性图快照,也不能自动观测硬删除。应用须在每次需要传递的变化发生时生成新序号;需要保留每次变化时,应写入独立的事件记录,而不是反复覆盖同一个节点。
配置
连接与数据库
neo4j.uri
配置 Neo4j 连接地址。
- 类型:
LIST - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:非空 URI 列表,多个地址以逗号分隔。支持
neo4j、neo4j+s、neo4j+ssc、bolt、bolt+s、bolt+ssc。首个 URI 用于创建驱动;配置多个 URI 时,全部地址均加入驱动的地址解析器。
neo4j.database
选择读取的数据库。
- 类型:
STRING - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:留空或空白使用驱动的默认数据库选择;修改数据库会改变源 Offset 身份。
认证
neo4j.authentication.type
选择数据库认证方式。
- 类型:
STRING - 默认值:
BASIC - 重要级别:高
- 有效值 / 注意事项:
NONE、BASIC、KERBEROS、BEARER、CUSTOM,区分大小写。对应方式的必需凭据不能为空。
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 - 重要级别:低
- 有效值 / 注意事项:仅接受
true、false。URI 使用+s或+ssc时由 URI 控制加密行为,此设置不覆盖它们。
neo4j.security.trust-strategy
选择显式加密连接的证书信任方式。
- 类型:
STRING - 默认值:
TRUST_SYSTEM_CA_SIGNED_CERTIFICATES - 重要级别:低
- 有效值 / 注意事项:
TRUST_ALL_CERTIFICATES、TRUST_SYSTEM_CA_SIGNED_CERTIFICATES、TRUST_CUSTOM_CA_SIGNED_CERTIFICATES。用于普通 URI 的显式加密,+s、+ssc使用自身信任行为;信任所有证书会削弱身份验证。
neo4j.security.hostname-verification-enabled
控制显式证书信任策略的主机名验证。
- 类型:
STRING - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:仅接受
true、false;使用显式信任策略时生效,不覆盖 URI 自带策略。
neo4j.security.cert-files
指定自定义 CA 证书文件。
- 类型:
LIST - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:使用自定义 CA 信任策略时提供 Worker 内可读的绝对文件路径,多个文件以逗号分隔;不能仅依赖界面提示判断是否需要证书。
连接池与重试
neo4j.connection-timeout
设置建立连接的超时时间。
- 类型:
STRING - 默认值:
30s - 重要级别:低
- 有效值 / 注意事项:非负整数配合小写单位
ms、s、m、h、d,可组合如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 - 重要级别:高
- 有效值 / 注意事项:
QUERY或CDC。QUERY 需要查询和 Topic;CDC 至少需要一个 Topic Pattern,且不接受RAW_JSON_STRING。
neo4j.start-from
选择没有可复用 Offset 时的起始位置。
- 类型:
STRING - 默认值:
NOW - 重要级别:高
- 有效值 / 注意事项:
EARLIEST、NOW、USER_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 - 重要级别:中
- 有效值 / 注意事项:仅接受
true、false。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 - 重要级别:低
- 有效值 / 注意事项:
true、false;仅适用于 QUERY,须与消费端预期的数据结构一致。
批量与负载格式
neo4j.batch-size
限制每批读取的数量。
- 类型:
INT - 默认值:
1000 - 重要级别:中
- 有效值 / 注意事项:至少为
1。QUERY 限制结果行,CDC 限制变化事件;CDC 路由展开后的消息数可能更大。不是字节限制,增大时需关注内存和下游压力。
neo4j.payload-mode
选择图数据的负载表示方式。
- 类型:
STRING - 默认值:
EXTENDED - 重要级别:中
- 有效值 / 注意事项:
EXTENDED、COMPACT、RAW_JSON_STRING。EXTENDED 使用类型包装,COMPACT 更紧凑但需关注类型变化;RAW_JSON_STRING 仅用于 QUERY。负载模式不替代 Converter,也不自动提供 Schema Registry 或 Schema 演进保障。
CDC 轮询与读取选项
neo4j.cdc.use-leader
配置 CDC 的 Leader 读取选项。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true、false。不要将此选项作为强制 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 默认值
- 重要级别:未声明
- 有效值 / 注意事项:
create、update、delete,运行时转为小写;须存在相同索引的 Pattern。省略时不设置操作过滤。
neo4j.cdc.topic.<topic>.patterns.<index>.changesTo
按发生变化的属性过滤事件。
- 类型:
STRING - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:逗号分隔的属性名,会去除两端空白;须存在相同索引的 Pattern。
neo4j.cdc.topic.<topic>.patterns.<index>.metadata.<metadata-key>
按事件元数据过滤该索引的变化。
- 类型:
STRING - 默认值:无 ConfigDef 默认值
- 重要级别:未声明
- 有效值 / 注意事项:metadata-key 支持
authenticatedUser、executingUser、txMetadata.<key>;键名支持字母、数字、点、下划线和连字符。须存在相同索引的 Pattern。
neo4j.cdc.topic.<topic>.key-strategy
选择 CDC 消息 Key 的构造方式。
- 类型:
STRING - 默认值:无 ConfigDef 默认值;未设置时运行时使用
WHOLE_VALUE - 重要级别:未声明
- 有效值 / 注意事项:
SKIP、ELEMENT_ID、ENTITY_KEYS、WHOLE_VALUE,区分大小写。ENTITY_KEYS 在实体没有键时输出 null;不要据此假定所有事件都有业务键。
neo4j.cdc.topic.<topic>.value-strategy
选择 CDC 消息 Value 的构造方式。
- 类型:
STRING - 默认值:无 ConfigDef 默认值;未设置时运行时使用
CHANGE_EVENT - 重要级别:未声明
- 有效值 / 注意事项:
CHANGE_EVENT或ENTITY_EVENT,区分大小写;前者保留变化事件表示,后者使用实体事件表示。
CDC 指标采集选项
neo4j.cdc.metric.last-db-tx-id.enabled
控制是否采集数据库最新事务 ID 指标。
- 类型:
STRING - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:仅接受
true、false;启用后会增加相应刷新查询。
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.query、neo4j.query.topic,覆盖 neo4j.start-from 并添加以下配置。将动态键中的 <cdc-topic> 替换为真实 Topic 名。
维护后从指定位置重放查询记录
适用业务场景:使用 QUERY 的流需要补发某个已知游标之后的数据,例如修复下游处理问题后重新处理记录;源端仍保留所需事件,消费端已准备好去重或幂等处理。 配置示例:在快速开始的 QUERY 配置上覆盖起始位置并添加以下配置。<last-processed-sequence> 为希望重放范围之前的最后一个序号;这里是源端数值位置,不是 Kafka Offset。
neo4j.ignore-stored-offset、neo4j.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负载模式。