> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Neo4j Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Neo4j Sink Connector，包括 Cypher 与 Pattern 策略、交付语义、安全、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=org.neo4j.connectors.kafka.sink.Neo4jConnector
topics=people
neo4j.uri=neo4j://<neo4j-host>:7687
neo4j.database=<neo4j-database>
neo4j.authentication.type=BASIC
neo4j.authentication.basic.username=<neo4j-username>
neo4j.authentication.basic.password=<neo4j-password>
neo4j.cypher.topic.people=MERGE (p:Person {id: __value.id}) SET p.name = __value.name, p.surname = __value.surname
```

替换 Neo4j 地址、数据库和凭证占位符。示例要求 `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` 和 `+ssc` URI 由 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`。

```properties theme={null}
connector.class=org.neo4j.connectors.kafka.sink.Neo4jConnector
topics=people
neo4j.uri=neo4j://<neo4j-host>:7687
neo4j.database=<neo4j-database>
neo4j.authentication.type=BASIC
neo4j.authentication.basic.username=<neo4j-username>
neo4j.authentication.basic.password=<neo4j-password>
neo4j.pattern.topic.people=(:Person{!id,name,surname})
```

关键说明：先确保消息值可转换为包含 `id`、`name` 和 `surname` 的 Map，并在目标数据库为业务唯一键建立合适的约束。Pattern 中的 `!id` 把 `id` 标记为节点匹配键，`name` 和 `surname` 作为写入属性；缺少有效唯一性时，可能匹配或创建非预期实体。不要让同一个 `people` Topic 同时配置 Cypher、CUD 或 CDC 策略。

### 持续写入时按分区扩展 Task

适用业务场景：Connector 已持续写入 Neo4j，单 Task 的吞吐或延迟不能满足要求，并且业务 Topic 有多个分区可供并行分配。以下配置在快速开始基础上把 Task 上限提高到 `4`。

```properties theme={null}
connector.class=org.neo4j.connectors.kafka.sink.Neo4jConnector
tasks.max=4
topics=people
neo4j.uri=neo4j://<neo4j-host>:7687
neo4j.database=<neo4j-database>
neo4j.authentication.type=BASIC
neo4j.authentication.basic.username=<neo4j-username>
neo4j.authentication.basic.password=<neo4j-password>
neo4j.cypher.topic.people=MERGE (p:Person {id: __value.id}) SET p.name = __value.name, p.surname = __value.surname
```

关键说明：`tasks.max=4` 只是请求的 Task 上限，不保证一定运行四个 Task；有效并行度不会超过可分配分区，也受 Neo4j 连接池和数据库写入容量限制。扩容前后应观察 Task 分配、吞吐、延迟和数据库负载。只依赖单个 Topic Partition 内的顺序，不要假定跨 Task 或跨分区存在全局顺序。

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

适用业务场景：持续写入期间可能出现少量无法转换、违反策略结构或导致 Neo4j 写入失败的记录，希望保留这些记录供排查，同时让可处理记录继续写入。以下配置启用错误容忍、DLQ 和上下文 Header。

```properties theme={null}
connector.class=org.neo4j.connectors.kafka.sink.Neo4jConnector
topics=people
neo4j.uri=neo4j://<neo4j-host>:7687
neo4j.database=<neo4j-database>
neo4j.authentication.type=BASIC
neo4j.authentication.basic.username=<neo4j-username>
neo4j.authentication.basic.password=<neo4j-password>
neo4j.cypher.topic.people=MERGE (p:Person {id: __value.id}) SET p.name = __value.name, p.surname = __value.surname
errors.tolerance=all
errors.deadletterqueue.topic.name=<dlq-topic-name>
errors.deadletterqueue.topic.replication.factor=3
errors.deadletterqueue.context.headers.enable=true
```

关键说明：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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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 或 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` 失败。
