概述
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。以下示例使用自管 Zeebe Gateway,record key 使用 String Converter,record value 使用 JSON Converter;record key 不参与 JSONPath 消息映射。如果使用 Camunda SaaS,请改用配置章节中的 SaaS 配置。<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。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 启动后没有发布消息?
检查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 的重试及恢复策略。