> ## 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.

# Camunda Zeebe Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Camunda Zeebe Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

Camunda Zeebe Sink Connector 从 Kafka 读取 JSON 业务事件，并将每条记录发布为 Zeebe 消息。它位于 Kafka 事件流和 Camunda Platform 8 流程之间：Connector 从事件中提取消息名称和关联键，向等待该消息的 BPMN 流程实例发起消息关联；也可以把事件中的变量和消息存活时间一并传给 Zeebe。

输入记录的 value 必须能按 JSON 和 JSONPath 语义解析。默认情况下，`messageName` 和 `correlationKey` 字段分别映射为 Zeebe 消息名称和关联键，`variables` 映射为流程变量，`timeToLive` 映射为消息 TTL。Kafka record key 不参与消息映射。Connector 可以连接 Camunda SaaS，也可以连接自管 Zeebe Gateway；它用于发布和关联消息，不用于启动流程、完成任务或消费 Zeebe 反馈。

## 前置条件

* Camunda SaaS 模式需要准备集群 ID、区域（非默认区域时）以及具有访问权限的 Client ID 和 Client Secret；自管模式需要准备可访问的 Zeebe Gateway 地址和对应认证方式。
* Kafka 输入 Topic 中的每条记录 value 必须是可解析的 JSON 文档，并能从配置的 JSONPath 读取 `messageName` 和 `correlationKey`；如果使用变量或 TTL，也要提供可转换为对应类型的字段。
* Zeebe 流程模型必须包含与消息名称和关联键匹配的消息接收事件；如果消息可能早于流程实例到达，应根据业务等待窗口设计 TTL。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka、可访问的 Zeebe Gateway 以及输入 Topic，确认网络连通和相应访问权限。创建和管理 Connector 的通用步骤参见 AutoMQ 的[管理 Connector](../manage-connectors)。以下示例使用自管 Zeebe Gateway，record key 使用 String Converter，record value 使用 JSON Converter；record key 不参与 JSONPath 消息映射。如果使用 Camunda SaaS，请改用配置章节中的 SaaS 配置。

```properties theme={null}
connector.class=io.zeebe.kafka.connect.ZeebeSinkConnector
tasks.max=1
topics=<kafka-topic>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
zeebe.client.gateway.address=<zeebe-gateway-host>:26500
zeebe.client.security.plaintext=false
message.path.messageName=$.messageName
message.path.correlationKey=$.correlationKey
message.path.variables=$.variables
message.path.timeToLive=$.timeToLive
```

将 `<kafka-topic>` 和 `<zeebe-gateway-host>` 替换为实际资源。每条输入记录的 value 至少应包含 `messageName` 和 `correlationKey`；`variables` 和 `timeToLive` 是可选字段。record key 可以是字符串，但不会替代 JSON value 中的字段，也不会参与 JSONPath 消息映射。自管网关使用明文传输时，将 `zeebe.client.security.plaintext` 设置为 `true`；Camunda SaaS 不使用这组 Gateway 配置，而是使用 `zeebe.client.cloud.*` 配置。应用配置后，Connector 会把输入记录发布为 Zeebe 消息，并由 Zeebe 按消息名称和关联键执行关联。

## 配置

### Zeebe Gateway

#### `zeebe.client.gateway.address`

自管 Zeebe Gateway 的地址。

* **类别**：Zeebe Gateway
* **类型**：`string`
* **默认值**：`localhost:26500`
* **重要级别**：高
* **有效值 / 注意事项**：填写可访问的 `host:port` 地址。仅在未设置 `zeebe.client.cloud.clusterId` 时生效；Camunda SaaS 模式会忽略此项。

#### `zeebe.client.requestTimeout`

Zeebe 请求超时时间，单位为毫秒。

* **类别**：Zeebe Gateway
* **类型**：`long`
* **默认值**：`1000`
* **重要级别**：低
* **有效值 / 注意事项**：直接 Gateway 模式使用此值；需要提供客户端可接受的非负毫秒时长。Camunda SaaS 客户端构造分支不使用此项。

#### `zeebe.client.security.plaintext`

是否使用明文连接自管 Zeebe Gateway。

* **类别**：Zeebe Gateway
* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：`true` 或 `false`。仅直接 Gateway 模式读取此项；只有在 Gateway 明确配置为明文传输时才使用 `true`，Camunda SaaS 模式会忽略此项。

### Camunda SaaS

#### `zeebe.client.cloud.clusterId`

要连接的 Camunda SaaS 集群 ID。

* **类别**：Camunda SaaS
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：设置非 `null` 值会选择 Camunda SaaS 客户端分支；应与 `zeebe.client.cloud.clientId` 和 `zeebe.client.cloud.clientSecret` 配合使用。该版本没有非空校验。

#### `zeebe.client.cloud.region`

Camunda SaaS 集群所在区域。

* **类别**：Camunda SaaS
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：只有同时设置 `zeebe.client.cloud.clusterId` 和此项时才覆盖 SDK 的区域选择；直接 Gateway 模式会忽略此项。区域值由 SDK 处理，本 Connector 不做区域校验。

### 身份验证

#### `zeebe.client.cloud.clientId`

连接 Camunda SaaS 或直接 OAuth 模式使用的 Client ID。

* **类别**：身份验证
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：Camunda SaaS 模式下与集群 ID、Client Secret 配合使用；直接 OAuth 模式下设置此项会启用 OAuth 凭据提供程序，并需要同时提供 Client Secret 和 token audience。

#### `zeebe.client.cloud.clientSecret`

连接 Camunda SaaS 或直接 OAuth 模式使用的 Client Secret。

* **类别**：身份验证
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：敏感值。通过受控的配置注入，不要写入共享文档或日志；Camunda SaaS 模式下与集群 ID、Client ID 配合使用。

#### `zeebe.client.cloud.token.audience`

直接 OAuth 模式请求令牌时使用的 token audience。

* **类别**：身份验证
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：仅在未设置 `zeebe.client.cloud.clusterId` 且设置了 `zeebe.client.cloud.clientId` 时使用；Camunda SaaS 客户端分支会忽略此项。该 Connector 不校验 audience 格式。

#### `zeebe.client.cloud.authorization.server.url`

注册表中用于表示授权服务器地址的配置项。

* **类别**：身份验证
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：在此版本中该配置项不会被客户端构造逻辑读取，设置它不会配置自定义 token endpoint。需要自定义 OAuth 端点时，不应依赖此项。

### 消息映射

#### `message.path.messageName`

从输入 JSON 文档提取 Zeebe 消息名称的 JSONPath 表达式。

* **类别**：消息映射
* **类型**：`string`
* **默认值**：`messageName`
* **重要级别**：高
* **有效值 / 注意事项**：填写可由 JSONPath 编译的表达式，并确保每条记录都能解析出可转换为消息名称的值。它是 JSONPath 表达式，不要求使用固定字段名。

#### `message.path.correlationKey`

从输入 JSON 文档提取 Zeebe 消息关联键的 JSONPath 表达式。

* **类别**：消息映射
* **类型**：`string`
* **默认值**：`correlationKey`
* **重要级别**：高
* **有效值 / 注意事项**：填写可由 JSONPath 编译的表达式，并确保每条记录都能解析出可转换为 Zeebe 关联键的值。建议为每个流程实例设计稳定且业务唯一的关联键。

#### `message.path.variables`

从输入 JSON 文档提取 Zeebe 消息变量的 JSONPath 表达式。

* **类别**：消息映射
* **类型**：`string`
* **默认值**：`variables`
* **重要级别**：中
* **有效值 / 注意事项**：空字符串可禁用变量路径。启用此项但记录中找不到路径时，Connector 会回退为将整个 JSON 文档作为变量。

#### `message.path.timeToLive`

从输入 JSON 文档提取 Zeebe 消息 TTL 的 JSONPath 表达式。

* **类别**：消息映射
* **类型**：`string`
* **默认值**：`timeToLive`
* **重要级别**：低
* **有效值 / 注意事项**：空字符串可禁用 TTL 路径；解析出的值按毫秒数读取。路径缺失时不发送 TTL，消息是否过期由 Zeebe 消息语义决定。

### Kafka Connect Sink 基础

#### `connector.class`

要实例化的 Kafka Connect Sink Connector 类。

* **类别**：Kafka Connect Sink 基础
* **类型**：`string`
* **默认值**：无 ConfigDef 默认值
* **重要级别**：高
* **有效值 / 注意事项**：必填，固定使用 `io.zeebe.kafka.connect.ZeebeSinkConnector`。

#### `tasks.max`

请求创建的最大 Sink Task 数量。

* **类别**：Kafka Connect Sink 基础
* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`。实际并行度还取决于输入 Topic 的分区分配和 Zeebe 客户端容量；增加该值不会提供跨 Task 的顺序保证。

#### `tasks.max.enforce`

是否强制执行 `tasks.max` 上限。

* **类别**：Kafka Connect Sink 基础
* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **已弃用**：是
* **有效值 / 注意事项**：`true` 或 `false`。Kafka Connect 已弃用此配置，未来主要版本可能移除；不建议新配置依赖 `false`。

#### `key.converter`

Connector 级别的 Kafka record key Converter。

* **类别**：Kafka Connect Sink 基础
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：必须是可实例化的 Converter 类；未设置时继承 Worker 的选择。推荐使用 `org.apache.kafka.connect.storage.StringConverter` 处理 record key；record key 仅完成 Kafka Connect 类型转换，不生成 Zeebe 消息字段，也不参与 JSONPath 消息映射。

#### `value.converter`

Connector 级别的 Kafka record value Converter。

* **类别**：Kafka Connect Sink 基础
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：必须是可实例化的 Converter 类；未设置时继承 Worker 的选择。转换后的 value 会按 JSON 语义解析，Connector 不自动选择 Converter。

#### `header.converter`

Connector 级别的 Kafka record header Converter。

* **类别**：Kafka Connect Sink 基础
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：必须是可实例化的 HeaderConverter 类；未设置时继承 Worker 的选择。本 Connector 不直接读取 record headers。

#### `transforms`

在 SinkTask 处理记录前执行的单消息转换（SMT）别名列表。

* **类别**：Kafka Connect Sink 基础
* **类型**：`list`
* **默认值**：空列表
* **重要级别**：低
* **有效值 / 注意事项**：使用逗号分隔的唯一别名；每个别名必须对应已定义的转换。SMT 先于 JSONPath 提取执行。

#### `predicates`

供条件 SMT 使用的谓词别名列表。

* **类别**：Kafka Connect Sink 基础
* **类型**：`list`
* **默认值**：空列表
* **重要级别**：低
* **有效值 / 注意事项**：使用逗号分隔的唯一别名；默认不启用任何谓词。

#### `config.action.reload`

外部 ConfigProvider 值发生变化时的处理动作。

* **类别**：Kafka Connect Sink 基础
* **类型**：`string`
* **默认值**：`restart`
* **重要级别**：低
* **有效值 / 注意事项**：有效值为 `none` 或 `restart`；它控制外部配置变化后的动作，不是 Zeebe 客户端属性。

### 错误处理

#### `errors.retry.timeout`

Kafka Connect 框架操作失败后的最长重试时长，单位为毫秒。

* **类别**：错误处理
* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：`0` 表示不重试，`-1` 表示无限重试。它与 Zeebe 请求超时以及 Connector 自身的异步重试机制不同。

#### `errors.retry.delay.max.ms`

Kafka Connect 框架重试之间的最大延迟，单位为毫秒。

* **类别**：错误处理
* **类型**：`long`
* **默认值**：`60000`
* **重要级别**：中
* **有效值 / 注意事项**：限制框架重试延迟；不会改变 Connector 内部 Zeebe 请求重试的退避参数。

#### `errors.tolerance`

Kafka Connect 对框架层错误的容忍策略。

* **类别**：错误处理
* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：有效值为 `none` 或 `all`。`all` 只允许 Connect 跳过框架处理阶段识别的问题，不能保证跳过 Zeebe 发布失败或所有 SinkTask 异常。

#### `errors.log.enable`

是否记录被容忍的错误和失败操作。

* **类别**：错误处理
* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：`true` 或 `false`。此项不会单独启用死信队列。

#### `errors.log.include.messages`

是否在错误日志中包含 Sink 记录的 Topic、分区、Offset 和时间戳。

* **类别**：错误处理
* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：`true` 或 `false`。开启后会增加记录元数据暴露范围；它不表示记录 value 会被写入日志。

#### `errors.deadletterqueue.topic.name`

Kafka Connect 框架错误报告器使用的死信队列 Topic。

* **类别**：错误处理
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：空字符串表示不配置 DLQ。配置后该 Topic 不能同时被 `topics` 消费或被 `topics.regex` 匹配；DLQ 不等于所有 Zeebe 发布失败都能被逐条恢复。

#### `errors.deadletterqueue.topic.replication.factor`

创建缺失 DLQ Topic 时使用的副本因子。

* **类别**：错误处理
* **类型**：`short`
* **默认值**：`3`
* **重要级别**：中
* **有效值 / 注意事项**：仅执行 `short` 类型解析；Broker 的 Topic 创建策略和副本约束仍然适用。

#### `errors.deadletterqueue.context.headers.enable`

是否在 DLQ 记录中写入 Kafka Connect 错误上下文 Header。

* **类别**：错误处理
* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：`true` 或 `false`。只有配置了 DLQ 目标时才生效，写入的 Header 使用 `__connect.errors.` 前缀。

### Kafka Topic 订阅

#### `topics`

Connector 要消费的 Kafka Topic 列表。

* **类别**：Kafka Topic 订阅
* **类型**：`list`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：使用逗号分隔的 Topic 名称；必须与 `topics.regex` 二选一，并且不能把配置的 DLQ Topic 纳入消费范围。

#### `topics.regex`

Connector 要消费的 Kafka Topic 正则表达式。

* **类别**：Kafka Topic 订阅
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：使用 Java Pattern 语法；必须与 `topics` 二选一，并且表达式不能匹配配置的 DLQ Topic。

## 最佳实践

### 按稳定关联键触发流程消息

适用业务场景：Kafka 持续接收订单、支付或库存事件，Zeebe 中的 BPMN 流程等待对应消息；需要让每个事件稳定地关联到目标流程实例，而不是依赖 Kafka record key。

```properties theme={null}
connector.class=io.zeebe.kafka.connect.ZeebeSinkConnector
tasks.max=1
topics=<business-events-topic>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
zeebe.client.gateway.address=<zeebe-gateway-host>:26500
zeebe.client.security.plaintext=false
message.path.messageName=$.messageName
message.path.correlationKey=$.orderId
message.path.variables=$.variables
```

关键说明：`correlationKey` 应映射到稳定且业务唯一的标识，例如订单 ID 或流程业务键，并与 BPMN 消息订阅的关联键设计一致。不要把 Connector 的重放行为当作跨系统 exactly-once；下游业务动作仍应具备幂等处理能力。如使用 SaaS，改用配置章节中的 `zeebe.client.cloud.*`。

### 用 TTL 缓冲先到达的事件

适用业务场景：业务事件可能先进入 Kafka，而等待消息的流程实例稍后才创建；需要让 Zeebe 在有限时间内保留消息，避免短暂的到达顺序差异直接丢失关联机会。

```properties theme={null}
connector.class=io.zeebe.kafka.connect.ZeebeSinkConnector
tasks.max=1
topics=<business-events-topic>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
zeebe.client.gateway.address=<zeebe-gateway-host>:26500
zeebe.client.security.plaintext=false
message.path.messageName=$.messageName
message.path.correlationKey=$.correlationKey
message.path.timeToLive=$.timeToLive
```

关键说明：`timeToLive` 的值按毫秒读取，只有输入文档能够通过配置路径提供该值时才会发送 TTL。TTL 到期表示 Zeebe 消息等待窗口结束，不表示 Kafka 记录被删除、业务流程已完成或消息进入 DLQ；TTL 应覆盖可接受的流程实例创建延迟。如使用 SaaS，改用配置章节中的 `zeebe.client.cloud.*`。

### 只传递流程需要的变量

适用业务场景：Kafka 事件包含较多业务字段，但流程只需要其中一组稳定、可序列化的字段；需要将这些字段作为 Zeebe message variables 传递，同时避免把无关数据带入流程上下文。

```properties theme={null}
connector.class=io.zeebe.kafka.connect.ZeebeSinkConnector
tasks.max=1
topics=<business-events-topic>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
zeebe.client.gateway.address=<zeebe-gateway-host>:26500
zeebe.client.security.plaintext=false
message.path.messageName=$.messageName
message.path.correlationKey=$.correlationKey
message.path.variables=$.workflowVariables
```

关键说明：将流程模型实际需要的字段放在 `workflowVariables` 对象中，并保持字段类型和结构稳定。如果配置的变量路径不存在，Connector 会回退为发送整个 JSON 文档；需要避免这种回退时，应在生产数据契约中强制提供该对象。该消息发布不提供流程完成回执，需要回执时另行设计 Zeebe 或 Kafka 反馈流。如使用 SaaS，改用配置章节中的 `zeebe.client.cloud.*`。

## 监控

### 监控内容

监控 Kafka Connect Worker、Connector 和 Task 的健康状态与重启次数，关注消费吞吐、处理延迟、Offset 提交进度、消费积压、错误和重试；同时观察 Worker JVM 的堆使用、垃圾回收、线程和 CPU。启用错误容忍或 DLQ 后，再分别监控被容忍错误、DLQ 写入量和 DLQ Topic 积压；Connector 内部 Zeebe 重试持续时，应结合 Task 状态、请求延迟和 Zeebe Gateway 可用性判断是否存在下游故障。

### 导入 Grafana 大盘

下载 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中选择已采集 Kafka Connect 指标的数据源，并确保指标标签包含集群、Connector 和 Task 标识，然后使用 Grafana 的 Import 功能导入该 JSON 大盘。

## 限制条件

* Connector 与 Kafka Offset 提交之间没有跨系统事务，不能保证 Kafka 到 Zeebe 的 exactly-once、无重复发布或原子批次。
* Connector 不提供同一 Kafka 分区、同一关联键或同一流程实例的严格消息顺序；`tasks.max` 增大后，各 Task 之间没有共享排序协调。
* 一个 Kafka Connect 批次中的 Zeebe 发布请求会并发执行；部分请求成功而其他请求失败时没有事务性回滚，重启或 Offset 提交间隙可能重新发布记录。
* Zeebe 客户端内部只对部分 gRPC 状态执行异步重试，且没有可配置的最大重试次数；JSONPath 解析错误和其他不可重试异常不会进入这组重试。
* `errors.tolerance` 和 DLQ 配置属于 Kafka Connect 框架错误处理，不能保证将 Connector 内部的 Zeebe 发布失败逐条写入 DLQ 或跳过。
* Connector 只按 JSON/JSONPath 语义读取消息，不应据此推断任意 Connect Schema、字段映射或 Schema 演进组合都兼容。
* 消息 TTL 到期后，Zeebe 可能丢弃仍未关联的消息；TTL 到期不是 Kafka 记录删除、流程完成或业务处理确认。

## 常见问题

### 为什么 Task 启动后没有发布消息？

检查 `topics` 和 `topics.regex` 是否恰好配置了一个订阅选择器，确认输入 Topic 有新记录，并检查每条 value 是否能按 `message.path.messageName` 和 `message.path.correlationKey` 解析。若使用自管 Zeebe，确认 Gateway 地址、TLS 或明文设置与服务端一致；若使用 SaaS，确认已设置集群 ID、Client ID 和 Client Secret，并且没有误用自管 Gateway 配置。

### 为什么消息没有关联到流程实例？

先确认 JSONPath 读取出的消息名称与 BPMN 消息订阅名称一致，再确认关联键的值与等待中的流程实例匹配。若事件可能早于流程实例到达，设置 `message.path.timeToLive` 并提供足够的毫秒值；TTL 到期后消息不会继续等待，也不会自动生成业务回执。

### 输入记录需要什么 JSON 结构？

默认结构至少包含 `messageName` 和 `correlationKey`，例如一个事件对象可以包含这两个字段以及 `variables` 和 `timeToLive`。如果字段位于其他路径，修改对应的 `message.path.*` 配置。value Converter 转换后的值仍必须符合 Connector 支持的 JSON/JSONPath 处理方式；Kafka record key 不会替代 JSON 中的关联键。

### 重启后为什么可能再次发布同一条消息？

Kafka Connect 只有在 `put` 成功并完成后续 Offset 提交流程时才推进记录位置；在 Zeebe 已接受消息但 Kafka Offset 尚未提交时发生重启、rebalance 或故障，记录可能被重新交付。Connector 会根据相同的 Topic、分区和 Offset 生成相同的消息 ID，Zeebe 对已存在的相同发布请求可按 `ALREADY_EXISTS` 视为成功，但这不是跨系统 exactly-once 保证，业务处理仍应设计为可重放或幂等。

### 为什么配置死信队列后仍看不到 Zeebe 发布失败记录？

DLQ 是 Kafka Connect 框架错误处理能力，主要覆盖 Converter、SMT 等框架阶段的错误；Connector 内部 Zeebe 发布失败会在 SinkTask 中处理，不能仅凭 `errors.deadletterqueue.topic.name` 保证逐条写入 DLQ。应同时检查 Task 错误日志、Zeebe Gateway 状态、认证配置和 Connector 的重试及恢复策略。
