Skip to main content

概述

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。以下示例使用未启用 TLS 和身份认证的本地 Zeebe gateway,且 gateway 与 Connect Worker 位于同一网络命名空间,可通过 localhost:26500 访问。 在 Zeebe 中部署包含消息发送服务任务的流程:Job 类型设为 kafka,自定义 header 的键设为 kafka-topic、值设为 workflow-events。提前准备该 Kafka Topic,并启动流程实例,使其到达这个服务任务。填写以下 Connector 配置:
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 提交超时。

最佳实践

首次接入时明确消息内容

适用业务场景:流程已有消息发送步骤,需要接入下游消费者。接入前明确消费者需要的变量范围和消息编码,避免默认将全部流程变量暴露给消费者。 配置示例:在快速开始配置上添加以下设置;orderIdstatus 是示例流程中的变量名,需替换为消费者实际需要、且流程已提供的变量。
关键说明:先确认流程类型与 Topic header,再用一个流程实例核对消息中的变量、Job key 和任务外层内容,最后接入消费者。选择变量不改变 Job JSON 外层结构;消费者应保留 Job key 供重复处理识别。若业务确实需要全部变量,保留空列表即可,不必为了使用本示例而缩小范围。Connector 会完成服务任务,不应与负责同一任务业务处理的其他 worker 无意竞争。

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

适用业务场景:Connector 已持续运行,但 Kafka 写入变慢、短暂网络中断或 Task 重启可能导致任务再次被领取。需要让正常处理有足够时间,并让消费者识别重复任务。 配置示例:在快速开始配置上添加以下设置;30000 毫秒是调整示例,需根据实际写入与完成延迟确定。
关键说明:锁定时长应覆盖正常处理延迟并留出余量;过短会增加重投机会,过长则推迟未完成任务的再次领取。消费者应以 Job key 为幂等标识处理同一任务的重复记录。重启后恢复依赖 Zeebe 的任务状态与超时机制,而不是从 Connect offset 定位历史记录。避免无意丢弃记录的 SMT 或容忍错误策略:被过滤或容忍丢弃的记录仍可能触发任务完成。完成命令异步发送且没有等待结果,Kafka 写入与任务完成不是原子操作,因此增加锁定时长和消费者去重都不能提供端到端无丢失或 exactly-once 保障。

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

适用业务场景:接入范围从一种任务扩展到多种已有的消息发送任务,单个 Task 依次激活不同类型已影响处理延迟。此时可将不同类型分配到不同 Task,避免增加无法承担这些任务的空闲 Task。 配置示例:在快速开始配置上添加以下设置;示例流程已分别使用 order-eventsshipment-events 两种 Job 类型,两个类型的任务都需配置有效的 Topic header。
关键说明:这里覆盖默认 kafka 类型,只有列出的类型会被领取。扩展前确认全部目标类型均在列表中,并观察 Task 分配、吞吐与延迟;重配置可能重新分组类型。任务数不应仅因吞吐需求而超过不同类型数,一个热门类型不会自动拆分到多个 Task。不要为了提高并行度人为复制列表项或改造业务任务类型。更多 Task 可能增加总激活量,需要同步观察资源与任务积压;不承诺类型间、Task 间或 Kafka 分区间的全局顺序。

监控

监控内容

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

导入 Grafana 大盘

下载共享的 Kafka Connect Grafana Dashboard,确认 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 的支持情况。