Skip to main content

概述

Camunda Zeebe Sink Connector 从 Kafka 读取 JSON 业务事件,并将每条记录发布为 Zeebe 消息。它位于 Kafka 事件流和 Camunda Platform 8 流程之间:Connector 从事件中提取消息名称和关联键,向等待该消息的 BPMN 流程实例发起消息关联;也可以把事件中的变量和消息存活时间一并传给 Zeebe。 输入记录的 value 必须能按 JSON 和 JSONPath 语义解析。默认情况下,messageNamecorrelationKey 字段分别映射为 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 读取 messageNamecorrelationKey;如果使用变量或 TTL,也要提供可转换为对应类型的字段。
  • Zeebe 流程模型必须包含与消息名称和关联键匹配的消息接收事件;如果消息可能早于流程实例到达,应根据业务等待窗口设计 TTL。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、可访问的 Zeebe Gateway 以及输入 Topic,确认网络连通和相应访问权限。创建和管理 Connector 的通用步骤参见 AutoMQ 的管理 Connector。以下示例使用自管 Zeebe Gateway,record key 使用 String Converter,record value 使用 JSON Converter;record key 不参与 JSONPath 消息映射。如果使用 Camunda SaaS,请改用配置章节中的 SaaS 配置。
<kafka-topic><zeebe-gateway-host> 替换为实际资源。每条输入记录的 value 至少应包含 messageNamecorrelationKeyvariablestimeToLive 是可选字段。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
  • 重要级别:低
  • 有效值 / 注意事项truefalse。仅直接 Gateway 模式读取此项;只有在 Gateway 明确配置为明文传输时才使用 true,Camunda SaaS 模式会忽略此项。

Camunda SaaS

zeebe.client.cloud.clusterId

要连接的 Camunda SaaS 集群 ID。
  • 类别:Camunda SaaS
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:设置非 null 值会选择 Camunda SaaS 客户端分支;应与 zeebe.client.cloud.clientIdzeebe.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
  • 重要级别:低
  • 已弃用:是
  • 有效值 / 注意事项truefalse。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
  • 重要级别:低
  • 有效值 / 注意事项:有效值为 nonerestart;它控制外部配置变化后的动作,不是 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
  • 重要级别:中
  • 有效值 / 注意事项:有效值为 noneallall 只允许 Connect 跳过框架处理阶段识别的问题,不能保证跳过 Zeebe 发布失败或所有 SinkTask 异常。

errors.log.enable

是否记录被容忍的错误和失败操作。
  • 类别:错误处理
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse。此项不会单独启用死信队列。

errors.log.include.messages

是否在错误日志中包含 Sink 记录的 Topic、分区、Offset 和时间戳。
  • 类别:错误处理
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse。开启后会增加记录元数据暴露范围;它不表示记录 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
  • 重要级别:中
  • 有效值 / 注意事项truefalse。只有配置了 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。
关键说明:correlationKey 应映射到稳定且业务唯一的标识,例如订单 ID 或流程业务键,并与 BPMN 消息订阅的关联键设计一致。不要把 Connector 的重放行为当作跨系统 exactly-once;下游业务动作仍应具备幂等处理能力。如使用 SaaS,改用配置章节中的 zeebe.client.cloud.*

用 TTL 缓冲先到达的事件

适用业务场景:业务事件可能先进入 Kafka,而等待消息的流程实例稍后才创建;需要让 Zeebe 在有限时间内保留消息,避免短暂的到达顺序差异直接丢失关联机会。
关键说明:timeToLive 的值按毫秒读取,只有输入文档能够通过配置路径提供该值时才会发送 TTL。TTL 到期表示 Zeebe 消息等待窗口结束,不表示 Kafka 记录被删除、业务流程已完成或消息进入 DLQ;TTL 应覆盖可接受的流程实例创建延迟。如使用 SaaS,改用配置章节中的 zeebe.client.cloud.*

只传递流程需要的变量

适用业务场景:Kafka 事件包含较多业务字段,但流程只需要其中一组稳定、可序列化的字段;需要将这些字段作为 Zeebe message variables 传递,同时避免把无关数据带入流程上下文。
关键说明:将流程模型实际需要的字段放在 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,在 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 启动后没有发布消息?

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

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

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

输入记录需要什么 JSON 结构?

默认结构至少包含 messageNamecorrelationKey,例如一个事件对象可以包含这两个字段以及 variablestimeToLive。如果字段位于其他路径,修改对应的 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 的重试及恢复策略。