> ## 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 Source Connector

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

## 概述

Camunda Zeebe Source Connector 将工作流中的 Zeebe Job 发布到 Kafka，使流程中的消息发送步骤能够连接事件处理、业务通知和数据集成链路。它按 Job 类型主动领取任务，以任务的自定义 header 指定目标 Topic，并将 Job key 和已激活任务的 JSON 内容作为消息键和值。

该 Connector 同时承担 Zeebe worker 的职责，会发送任务完成命令，使流程能够继续执行。它不是被动的流程日志订阅器，也不导出所有流程事件或历史记录。一个有效任务对应一条发往单个 Topic 的记录；与其他 worker 使用相同 Job 类型时，各 worker 竞争领取任务，而不是各自获得一份广播副本。

## 前置条件

* Zeebe 流程中需要有可被激活的服务任务，其 Job 类型与 Connector 的 `job.types` 一致；任务的自定义 header 需要提供单个有效的 Kafka Topic 名称，并允许 Connector 完成该任务、推进流程。
* Self-Managed Zeebe 的 gRPC gateway 通常使用 `26500` 端口；连接时需匹配 gateway 的 TLS 与身份认证方式。使用 Camunda SaaS 时，需要集群 ID、区域及具备相应访问权限的客户端凭证。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和本地 Self-Managed Zeebe 8.4，确认访问权限与网络连通；Connector 的创建与管理操作参见[管理 Connector](../manage-connectors)。以下示例使用未启用 TLS 和身份认证的本地 Zeebe gateway，且 gateway 与 Connect Worker 位于同一网络命名空间，可通过 `localhost:26500` 访问。

在 Zeebe 中部署包含消息发送服务任务的流程：Job 类型设为 `kafka`，自定义 header 的键设为 `kafka-topic`、值设为 `workflow-events`。提前准备该 Kafka Topic，并启动流程实例，使其到达这个服务任务。填写以下 Connector 配置：

```properties theme={null}
connector.class=io.zeebe.kafka.connect.ZeebeSourceConnector
zeebe.client.security.plaintext=true
key.converter=org.apache.kafka.connect.converters.LongConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
```

gateway 默认地址为 `localhost:26500`，默认 Job 类型为 `kafka`，默认路由 header 为 `kafka-topic`，因此示例不重复设置这些值。如果 Worker 在独立容器或远程主机中运行，需将 `zeebe.client.gateway.address` 设为 Worker 实际可访问的 gateway 地址；容器中的 `localhost` 不能代指宿主机或另一个容器。明文连接仅适用于这里明确关闭 TLS 的本地环境。

配置生效后，检查 `workflow-events` 中是否出现消息，并在 Zeebe 中确认任务完成及流程继续执行。消息键是 LongConverter 编码的 64 位 Job key，不是文本数字；消息值是 StringConverter 编码的已激活 Job JSON，而不只是流程变量。消费者应解析任务外层内容后再读取变量，不要预设未经确认的变量字段表示方式。Kafka 中出现消息与 Zeebe 任务完成应分别检查：完成命令异步发送，不能只凭 Kafka 写入成功认定流程已经推进。

## 配置

### Connector 与输出编码

#### `connector.class`

选择 Zeebe Source Connector 实现。

* **类型**：`string`
* **默认值**：无 ConfigDef 默认值
* **重要级别**：高
* **必填**：是
* **有效值 / 注意事项**：使用 `io.zeebe.kafka.connect.ZeebeSourceConnector`，并确保该类可加载。

#### `tasks.max`

设置 Kafka Connect 可创建的最大 Task 数。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`。Connector 按 `job.types` 列表分配类型，一个唯一的类型不会自动拆到多个 Task。Task 数超过类型列表项数不会带来有用的并行度；不要用重复类型充当分片。

#### `key.converter`

覆盖 Worker 的消息键转换器。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未设置时继承 Worker 配置。键的 Connect 类型是 `INT64`，内容是 Zeebe Job key；快速开始显式使用 `org.apache.kafka.connect.converters.LongConverter`，这不是插件默认值。

#### `value.converter`

覆盖 Worker 的消息值转换器。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未设置时继承 Worker 配置。值的 Connect 类型是 `STRING`，内容为已激活 Job 的 JSON 文本。`org.apache.kafka.connect.storage.StringConverter` 可直接编码该文本；JsonConverter 处理的是字符串值，不会自动将其转换成结构化 Connect 对象。

### Self-Managed 连接与请求

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

指定直接连接的 Zeebe gateway 地址。

* **类型**：`string`
* **默认值**：`localhost:26500`
* **重要级别**：高
* **有效值 / 注意事项**：使用 gateway 的主机与端口。仅在未设置 `zeebe.client.cloud.clusterId` 时生效；云连接模式不使用此项。

#### `zeebe.client.requestTimeout`

设置直接连接客户端的默认请求超时，同时作为 Source 轮询退避时长的上界，单位为毫秒。

* **类型**：`long`
* **默认值**：`1000`
* **重要级别**：低
* **有效值 / 注意事项**：ConfigDef 未声明范围校验；应使用可用的非负时长。激活请求会使用当前退避时长覆盖客户端默认请求超时，因此此项不是每次激活请求的固定超时。云模式不将此项应用为客户端默认请求超时，但 Source 退避仍读取它。

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

控制直接连接是否关闭传输层 TLS。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：仅在明确使用明文 gateway 时设置为 `true`。云模式忽略此项；不要为绕过 TLS 配置问题而关闭安全连接。

### 云集群选择

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

选择 Camunda SaaS 连接模式与集群。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：任何非 `null` 值都会选择云模式，包括空字符串。使用 Self-Managed 时应移除此项，而不是填空字符串。云模式需提供相应客户端凭证；Connector 不自行检查凭证是否完整，且不使用直接连接的 gateway、明文开关及 OAuth audience 设置。

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

指定云集群区域。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：仅在云模式且本项非 `null` 时传给客户端。应与集群实际区域一致；未设置时由 SDK 处理，不代表 Connector 提供了固定区域默认值。

### 身份认证

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

指定 OAuth 客户端 ID。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：云模式使用此 ID；直接连接模式中，非 `null` 的 ID 会启用 OAuth 凭证提供器，并使用 secret 与 audience。Connector 不自行校验凭证配对，需按 gateway 的认证要求准备。

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

指定 OAuth 客户端密钥。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：云模式使用此值；直接连接仅在设置客户端 ID 时使用。它注册为 `STRING` 而不是 `PASSWORD`，不要假设配置展示会自动脱敏；按平台的密钥管理方式提供，不在日志或共享配置中暴露。属性名是 `clientSecret`，不要替换成 SDK 的其他属性名。

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

指定直接连接 OAuth 凭证的 token audience。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：仅在未设置云集群 ID、且已设置客户端 ID 时使用。填写认证系统要求的 audience；云模式忽略此项，`null` 不是某个固定 audience 字符串。

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

指定授权服务器 URL，但本版本未使用该设置。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：在 Connector 0.51.0 中未应用到客户端构建，修改它不能设置 token 获取端点。不要将其作为可用的自定义 OAuth 授权 URL 配置。

### Job 选择与路由

#### `job.types`

选择需要激活的 Zeebe Job 类型。

* **类型**：`list`
* **默认值**：`kafka`
* **重要级别**：高
* **有效值 / 注意事项**：逗号分隔的类型列表，解析后不能为空。匹配服务任务的 Job 类型，不支持正则表达式，也不是流程事件类型过滤器。列表不会自动去重；使用互不重复的类型。

#### `job.header.topics`

指定包含目标 Kafka Topic 名称的 Zeebe 自定义 header 键。

* **类型**：`string`
* **默认值**：`kafka-topic`
* **重要级别**：高
* **有效值 / 注意事项**：配置值是 header 的键名，不是 Topic 名称。每个 Job 必须在该键下提供一个非空、有效的单个 Topic 名称；header 的完整值直接用作目标 Topic，不按逗号拆分，也不执行多 Topic fan-out。仅非空检查不能保证 Kafka 名称合法。

#### `job.variables`

选择激活任务时获取的流程变量。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：低
* **有效值 / 注意事项**：使用逗号分隔的变量名；空列表表示获取全部变量。此项仅缩小获取的变量范围，不筛选 Job，也不去掉任务外层元数据与自定义 header。

### Job 激活与超时

#### `zeebe.client.worker.name`

指定领取任务时使用的 Zeebe worker 名称。

* **类型**：`string`
* **默认值**：`kafka-connector`
* **重要级别**：低
* **有效值 / 注意事项**：原样传给任务激活命令，不是 Kafka Connect Worker 名称，也不形成独占锁或任务分片。

#### `zeebe.client.worker.maxJobsActive`

设置每个 Task、每个 Job 类型的激活与在途容量目标。

* **类型**：`int`
* **默认值**：`100`
* **重要级别**：中
* **有效值 / 注意事项**：ConfigDef 未声明范围校验；应使用正整数，零或负数没有可激活容量。它不是整个 Connector 的严格在途上限，也不是一次轮询的全局批量大小；调整时需结合写入延迟、积压和 Worker 资源观察实际运行情况。

#### `zeebe.client.job.timeout`

设置 Job 激活后的锁定时长，单位为毫秒。

* **类型**：`long`
* **默认值**：`5000`
* **重要级别**：中
* **有效值 / 注意事项**：ConfigDef 未声明范围校验；应使用能覆盖实际 Kafka 写入及 Zeebe 完成延迟的正时长。超时后未完成任务可被再次激活，Connector 不自动续期。它不是激活 RPC 请求超时，也不是 Connect offset 提交超时。

## 最佳实践

### 首次接入时明确消息内容

**适用业务场景**：流程已有消息发送步骤，需要接入下游消费者。接入前明确消费者需要的变量范围和消息编码，避免默认将全部流程变量暴露给消费者。

**配置示例**：在快速开始配置上添加以下设置；`orderId` 与 `status` 是示例流程中的变量名，需替换为消费者实际需要、且流程已提供的变量。

```properties theme={null}
job.variables=orderId,status
```

**关键说明**：先确认流程类型与 Topic header，再用一个流程实例核对消息中的变量、Job key 和任务外层内容，最后接入消费者。选择变量不改变 Job JSON 外层结构；消费者应保留 Job key 供重复处理识别。若业务确实需要全部变量，保留空列表即可，不必为了使用本示例而缩小范围。Connector 会完成服务任务，不应与负责同一任务业务处理的其他 worker 无意竞争。

### 日常运行中为故障恢复留出时间

**适用业务场景**：Connector 已持续运行，但 Kafka 写入变慢、短暂网络中断或 Task 重启可能导致任务再次被领取。需要让正常处理有足够时间，并让消费者识别重复任务。

**配置示例**：在快速开始配置上添加以下设置；`30000` 毫秒是调整示例，需根据实际写入与完成延迟确定。

```properties theme={null}
zeebe.client.job.timeout=30000
```

**关键说明**：锁定时长应覆盖正常处理延迟并留出余量；过短会增加重投机会，过长则推迟未完成任务的再次领取。消费者应以 Job key 为幂等标识处理同一任务的重复记录。重启后恢复依赖 Zeebe 的任务状态与超时机制，而不是从 Connect offset 定位历史记录。避免无意丢弃记录的 SMT 或容忍错误策略：被过滤或容忍丢弃的记录仍可能触发任务完成。完成命令异步发送且没有等待结果，Kafka 写入与任务完成不是原子操作，因此增加锁定时长和消费者去重都不能提供端到端无丢失或 exactly-once 保障。

### 按已有任务类型增加并行处理

**适用业务场景**：接入范围从一种任务扩展到多种已有的消息发送任务，单个 Task 依次激活不同类型已影响处理延迟。此时可将不同类型分配到不同 Task，避免增加无法承担这些任务的空闲 Task。

**配置示例**：在快速开始配置上添加以下设置；示例流程已分别使用 `order-events` 和 `shipment-events` 两种 Job 类型，两个类型的任务都需配置有效的 Topic header。

```properties theme={null}
job.types=order-events,shipment-events
tasks.max=2
```

**关键说明**：这里覆盖默认 `kafka` 类型，只有列出的类型会被领取。扩展前确认全部目标类型均在列表中，并观察 Task 分配、吞吐与延迟；重配置可能重新分组类型。任务数不应仅因吞吐需求而超过不同类型数，一个热门类型不会自动拆分到多个 Task。不要为了提高并行度人为复制列表项或改造业务任务类型。更多 Task 可能增加总激活量，需要同步观察资源与任务积压；不承诺类型间、Task 间或 Kafka 分区间的全局顺序。

## 监控

### 监控内容

关注 Kafka Connect 的 Connector / Task 状态、吞吐与端到端延迟、Offset 提交耗时及失败、错误与重试，以及 Worker JVM 的堆内存、GC 和线程信号；Task 处于 RUNNING 不等于已经成功领取或发布任务，需要结合实际消息与日志判断。仅在部署启用了相应错误处理时关注 DLQ 活动，不将 DLQ 活动视为 Zeebe 任务完成的证明。

### 导入 Grafana 大盘

下载共享的 [Kafka Connect Grafana Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认 Kafka Connect 指标已采集到兼容的 Prometheus 数据源，且采集标签与大盘使用的集群、Connector 和 Task 标签匹配；随后在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* 数据入口是指定类型的可激活 Job，不支持全量流程事件导出、历史记录回放，或按流程 ID、元素、租户及任意条件筛选任务。
* 每个 Job 只路由到一个 Topic；自定义 header 中的逗号分隔字符串不会展开为多个目标 Topic。
* 一个唯一 Job 类型不会由该 Connector 自动拆分到多个 Task，增加 `tasks.max` 不会直接提升单类型的并行度。
* Kafka 写入与 Zeebe 任务完成不构成跨系统原子事务；完成命令未等待异步确认，不提供端到端 exactly-once 或无条件的至少一次交付保障。
* Connect offset 包含 Job key，但不用于启动时定位或过滤任务；重置 offset 不会重新发布已经完成的 Job。
* Connector 不自动续期激活锁，也不保证停止时所有异步完成命令均已确认。
* `zeebe.client.cloud.authorization.server.url` 在 Connector 0.51.0 中不生效，不能用于设置自定义 token 获取端点。

## 常见问题

### Task 显示 RUNNING，为什么 Topic 没有消息？

检查流程是否已到达可激活的服务任务，Job 类型是否与 `job.types` 一致，是否有其他 worker 竞争相同类型。再检查 gateway 地址、TLS 和认证模式以及激活错误日志；持续的授权错误也可能表现为空轮询。确认 header 指向单个合法 Topic，Topic 已存在且允许写入，并检查 Converter 与 Kafka 写入错误。不要仅凭运行状态认定连通和发布成功。

### 为什么缺少 Topic header 后任务失败或出现 incident？

任务缺少指定 header 或其值为空时，Connector 不生成输出记录，而会异步发送任务失败命令，将剩余重试数减少一。若失败命令被 Zeebe 接受，仍有重试次数的任务可再次领取，次数耗尽则产生 incident。检查 `job.header.topics` 所指定的键与流程定义是否一致，修正流程中的 Topic header；对已有失败任务，按 Zeebe 的 incident 与重试管理方式处理。非空但非法的 Topic 名称可能通过本地检查，却在 Kafka 写入阶段失败，应检查 Kafka 错误，而不是期待同一缺失 header 路径处理它。

### 为什么设置两个 Topic 名称没有得到两份消息？

header 值不会按逗号拆分，而是整体当作一个 Topic 名称。将 header 改为单个有效 Topic；如业务需要多个目的地，应在工作流中设计独立的消息发送步骤或在 Kafka 下游实现分发，不依赖此 Connector 的 header 实现 fan-out。

### 为什么消息中有任务元数据，而不只有变量？

消息值是已激活 Job 的 JSON，包含任务外层内容。`job.variables` 只控制获取哪些变量，不剥离外层结构。按实际编码解析消息并提取所需变量；若使用 JsonConverter，还需区分它对字符串的包装与 Job JSON 本身。需要直接读取 JSON 文本时，可沿用快速开始的 StringConverter。

### 为什么 Kafka 已有消息，流程却没有继续？

任务完成命令异步发送，Connector 不等待或检查异步结果，Kafka 写入成功不证明 Zeebe 已确认完成。检查任务在 Zeebe 中的实际状态、gateway 连接和认证、是否已超时并被重新激活。恢复连接并处理仍未完成的任务，按业务要求核对已发布消息与流程状态；不要以重置 Connect offset 作为流程修复手段。

### 为什么重启后出现相同 Job key，重置 offset 却不能回放？

Kafka 写入后、Zeebe 完成确认前发生故障或锁定超时，未完成任务可能再次激活并使用相同 Job key。消费者应按该键实现幂等处理，并根据实际延迟调整 Job 锁定时长。Connect offset 不作为可回放游标读取；已经完成的 Job 不会因 offset 重置重新激活。需要重发业务消息时，应通过业务补偿或新的流程任务处理。

### 为什么修改授权服务器 URL 仍未改变认证端点？

Connector 0.51.0 虽注册了 `zeebe.client.cloud.authorization.server.url`，但未将它应用到客户端。确认使用的是云模式还是直接连接模式，并核对适用的客户端 ID、secret 与 audience；不能通过反复修改此项切换 token 端点。若部署要求自定义授权端点，需要选用明确支持该认证方式的集成方案，或在升级前确认目标 Connector 的支持情况。
