概述
Amazon EventBridge Sink Connector 将一个或多个 Kafka Topic 中的记录发送到指定的 Amazon EventBridge 事件总线。每条 Kafka 记录会转换为一个 EventBridge 事件,Kafka 的 Topic、分区、Offset、时间戳、Header、Key 和 Value 保存在事件的detail 中,事件的 source 使用 Connector 标识生成。它适合把 Kafka 中的业务事件、CDC 记录或应用消息接入 EventBridge,再由 EventBridge 规则路由到下游 AWS 服务或其他目标。
默认情况下,事件的 detail-type 根据 Topic 生成,事件时间由 EventBridge 在接收时设置。可以按 Topic 或消息内容自定义 detail-type,也可以将事件 Value 的指定部分外置到 S3,并在 EventBridge 事件中保留引用信息。
前置条件
- 目标 EventBridge 事件总线已创建,Connector 使用的 AWS 身份具有该事件总线的
events:PutEvents权限;启用 IAM 角色、配置文件或自定义凭据提供程序时,还需准备相应的信任关系和凭据来源。 - 如果启用 S3 外置,目标 S3 Bucket 已创建,Connector 使用的 AWS 身份具有该 Bucket 的
s3:PutObject权限。 - 如果使用 Avro、Protobuf、AWS Glue Schema Registry 或自定义 Converter,相关 Converter 类及其依赖已放入 Kafka Connect Worker 可加载的插件路径,并准备好读取 Schema 所需的权限。
- 如果配置自定义 detail-type、时间或凭据提供程序类,相关类及其依赖已放入 Kafka Connect Worker 的插件路径,并满足对应接口和无参构造要求。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、输入 Topic 和目标 EventBridge 事件总线,并确认网络连通及访问权限;创建和管理操作参见管理 Connector。以下配置使用无 Schema JSON 值,将记录发送到指定事件总线。detail-type 默认使用 kafka-connect-<topic>,detail 包含 Kafka 记录元数据及解码后的 Key、Value。value.converter.schemas.enable=false 适用于不带 Schema 外壳的 JSON 值;使用其他消息格式时,改用对应的 Converter 及其依赖。
配置
输入与任务
topics
指定 Sink Connector 消费的 Kafka Topic 列表。
- 类型:
list - 默认值:空列表
[] - 重要级别:高
- 有效值 / 注意事项:使用逗号分隔的 Topic 名称;与
topics.regex二选一,至少配置一个非空选项。
topics.regex
使用正则表达式选择 Sink Connector 消费的 Kafka Topic。
- 类型:
string - 默认值:空字符串
"" - 重要级别:高
- 有效值 / 注意事项:使用 Java 正则表达式;与
topics互斥,不能同时配置。
tasks.max
请求为该 Connector 创建的最大 Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1。实际并行度还受 Kafka 分区、Worker 资源和 EventBridge 吞吐限制;跨 Task、Topic 或分区不提供全局顺序保证。
EventBridge 连接与目标
aws.eventbridge.connector.id
Connector 的唯一标识,用于生成 EventBridge 事件的 source,也用于 IAM 角色会话名称。
- 类型:
string - 默认值:无固定默认值,必填
- 重要级别:高
- 有效值 / 注意事项:不能为空或只包含空格;没有额外的字符或长度限制。事件的
source形如kafka-connect.<connector-id>。
aws.eventbridge.region
目标 EventBridge 事件总线所在的 AWS 区域。
- 类型:
string - 默认值:无固定默认值,必填
- 重要级别:高
- 有效值 / 注意事项:使用 AWS SDK 支持的区域名称,例如
us-east-1。
aws.eventbridge.eventbus.arn
目标 EventBridge 事件总线的 ARN。
- 类型:
string - 默认值:无固定默认值,必填
- 重要级别:高
- 有效值 / 注意事项:使用目标事件总线的完整 ARN,例如
arn:aws:events:us-east-1:123456789012:event-bus/orders。
aws.eventbridge.endpoint.uri
覆盖 AWS SDK 默认的 EventBridge 服务端点。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空使用 AWS SDK 默认端点;非空值必须是可解析的 URI,错误 URI 可能导致 Task 启动失败。
aws.eventbridge.eventbus.global.endpoint.id
指定 EventBridge 全局端点 ID。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空表示不发送全局端点 ID;非空值使用两段式格式,例如
abcde.veo。
AWS 身份认证
aws.eventbridge.auth.credentials_provider.class
指定自定义 AWS 凭据提供程序类。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空使用默认凭据提供程序;非空类必须实现
AwsCredentialsProvider,并具有无参构造方法。实现Configurable时,Connector 的原始配置会传给该类。
aws.eventbridge.iam.role.arn
指定 Connector 使用 STS 承担的 IAM 角色 ARN。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空不启用该角色;非空值使用 IAM 角色 ARN 格式,例如
arn:aws:iam::123456789012:role/EventBridgePutEventsRole。未配置自定义凭据提供程序时,Connector 使用 STS 承担该角色。
aws.eventbridge.iam.external.id
为基于 IAM 角色的认证提供 External ID。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:仅在使用
aws.eventbridge.iam.role.arn时生效。External ID 属于敏感配置,不要写入日志、示例或提交到公共仓库。
aws.eventbridge.iam.profile.name
指定 AWS 共享配置文件中的凭据 Profile。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:仅适用于默认凭据提供程序路径。使用该配置时不要同时设置
AWS_PROFILE、AWS_ACCESS_KEY_ID、AWS_SECRET_ACCESS_KEY或AWS_SESSION_TOKEN,否则会与配置文件选择冲突。
事件映射
aws.eventbridge.detail.types
为 EventBridge 事件设置 detail-type,可按 Topic 映射或对所有 Topic 使用同一表达式。
- 类型:
list - 默认值:
[kafka-connect-${topic}] - 重要级别:中
- 有效值 / 注意事项:可使用单个表达式,例如
orders-${topic};也可使用静态值,例如business-event;多个值时使用topic:detail-type,例如orders:order-created,customers:customer-updated。未命中的 Topic 回退到kafka-connect-<topic>。配置自定义 detail-type mapper 后,该项不再生效。
aws.eventbridge.detail.types.mapper.class
指定从 Topic 或记录内容计算 detail-type 的类。
- 类型:
string - 默认值:
software.amazon.event.kafkaconnector.mapping.DefaultDetailTypeMapper - 重要级别:中
- 有效值 / 注意事项:类必须实现
DetailTypeMapper并具有无参构造方法。使用内置JsonPathDetailTypeMapper时,必须同时配置aws.eventbridge.detail.types.jsonpathmapper.fieldref;配置该类后,aws.eventbridge.detail.types会被忽略。
aws.eventbridge.detail.types.jsonpathmapper.fieldref
指定 JsonPathDetailTypeMapper 从 Kafka 记录 Value 中提取 detail-type 的 JSONPath。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:使用可解析且 definite 的 JSONPath,例如
$.metadata.event-type;提取结果必须是非空字符串。路径不存在、结果为空、为null或不是字符串时回退到记录 Topic。该配置仅在使用JsonPathDetailTypeMapper时需要。
aws.eventbridge.time.mapper.class
指定为每条事件计算 EventBridge time 字段的类。
- 类型:
string - 默认值:
software.amazon.event.kafkaconnector.mapping.DefaultTimeMapper - 重要级别:中
- 有效值 / 注意事项:类必须可加载、实现
TimeMapper并具有无参构造方法。默认值不提供自定义时间,由 EventBridge 在PutEvents调用时设置时间。
aws.eventbridge.eventbus.resources
为每条 EventBridge 事件添加 resources 列表。
- 类型:
list - 默认值:空列表
[] - 重要级别:中
- 有效值 / 注意事项:使用逗号分隔的资源值;资源格式由 EventBridge 服务校验。
S3 载荷外置
aws.eventbridge.offloading.default.s3.bucket
指定用于外置事件载荷的 S3 Bucket。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空不启用外置;非空值会启用 claim-check 处理,并且 Bucket 必须已存在且允许当前 AWS 身份写入对象。
aws.eventbridge.offloading.default.fieldref
指定从事件 detail 中选择并外置到 S3 的 JSONPath。
- 类型:
string - 默认值:
$.detail.value - 重要级别:中
- 有效值 / 注意事项:仅在配置 S3 Bucket 时生效,路径必须是 definite JSONPath 且以
$.detail.value开头;不支持数组或通配符路径。匹配不到值或值为null时透传原事件,空对象和空数组会被外置。
aws.eventbridge.offloading.default.s3.endpoint.uri
覆盖 S3 外置使用的服务端点。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:仅在启用 S3 外置时使用;留空使用 AWS SDK 默认 S3 端点,非空值必须是可解析的 URI。
投递与错误处理
aws.eventbridge.retries.max
设置 EventBridge 发送失败时的 Connector 层最大重试次数。
- 类型:
int - 默认值:
2 - 重要级别:中
- 有效值 / 注意事项:取值范围为
0到10;0表示不进行 Connector 层重试。该值同时用于 AWS SDK 客户端,实际请求次数可能高于 Connector 层配置的次数。
aws.eventbridge.retries.delay
设置 Connector 两次重试之间的等待时间,单位为毫秒。
- 类型:
int - 默认值:
200 - 重要级别:中
- 有效值 / 注意事项:连接器不会拒绝负值;建议使用非负整数,因为负值不会产生有效等待时间。该项不改变错误是否可重试。
errors.tolerance
设置 Kafka Connect 对转换、转换器和 Task 错误的容忍模式。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可选
none或all。all通常应与 DLQ 配置一起使用;它与aws.eventbridge.retries.max的 EventBridge 发送重试不同。
errors.deadletterqueue.topic.name
指定接收由 Kafka Connect 错误报告器处理的失败记录的 DLQ Topic。
- 类型:
string - 默认值:空字符串
"" - 重要级别:中
- 有效值 / 注意事项:留空表示不配置 DLQ;配置后,Connector 可将记录级别的不可重试失败报告到该 Topic。Topic 的创建和权限由 Kafka Connect 与 Kafka 集群管理。
errors.deadletterqueue.topic.replication.factor
设置 Kafka Connect 创建 DLQ Topic 时使用的副本因子。
- 类型:
short - 默认值:
3 - 重要级别:中
- 有效值 / 注意事项:仅在配置
errors.deadletterqueue.topic.name时生效;取值必须与 Kafka 集群可用 Broker 数和创建策略匹配。
消息转换
key.converter
指定将 Kafka 消息 Key 解码为 Kafka Connect 值的 Converter。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时使用 Worker 的 Converter 配置;显式设置时必须是可实例化的 Converter 类。解码后的 Key 会写入 EventBridge 事件
detail.key。
value.converter
指定将 Kafka 消息 Value 解码为 Kafka Connect 值的 Converter。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:未设置时使用 Worker 的 Converter 配置;显式设置时必须是可实例化的 Converter 类。Avro、Protobuf 和 Schema Registry 依赖由 Worker 环境提供;解码失败的记录不会进入 EventBridge 发送批次。
最佳实践
将多个 Topic 映射为可路由的事件类型
适用业务场景:多个业务 Topic 需要进入同一个 EventBridge 事件总线,并希望通过稳定的detail-type 区分订单、客户等事件类型,以便配置规则时避免只依赖 Topic 名称。
配置示例:在快速开始配置中添加以下设置,保留原有事件总线和身份配置。
detail-type,未命中的 Topic 会回退到默认值。多个映射项必须使用 topic:detail-type 格式;如果改用自定义 mapper,应移除或忽略 aws.eventbridge.detail.types,避免误以为两种映射会叠加。
将大载荷外置到 S3
适用业务场景:Kafka 记录包含较大的业务 Value,而 EventBridge 事件需要保留可过滤的元数据,同时让下游消费者按需读取完整载荷。 配置示例:在快速开始配置中添加以下设置,Bucket 使用已创建且允许写入的 S3 资源。dataref 和 datarefJsonPath 引用信息。目标 AWS 身份必须具有 s3:PutObject 权限;路径只能从 $.detail.value 开始,匹配不到或值为 null 时不外置。S3 外置不能替代下游对对象生命周期、访问权限和清理策略的管理。
配置失败重试并将不可恢复记录发送到 DLQ
适用业务场景:EventBridge 或网络的暂时性失败需要自动重试,同时转换失败或不可重试的记录不能阻塞后续消息,希望保留这些记录供排查和补偿。 配置示例:在快速开始配置中添加以下设置,并为 DLQ Topic 配置独立的消费和保留策略。aws.eventbridge.retries.max 控制 Connector 层的发送重试,DLQ 由 Kafka Connect 错误报告器处理,两者解决不同问题。AWS SDK 也会使用该重试值,因此实际调用次数可能超过 Connector 层的直观估算。上线前应为 DLQ 建立告警、保留和重放流程,不要把 errors.tolerance=all 当作成功投递保证。
监控
监控内容
监控 Kafka Connect Worker、Connector 和 Task 是否处于预期状态,Task 的记录吞吐、处理延迟、Consumer Lag、Offset 提交、失败记录、重试次数和错误日志;同时关注 Worker JVM 的堆内存、GC、线程和 CPU。启用errors.tolerance=all 与 DLQ 后,再监控 DLQ Topic 的写入量和积压,并将 DLQ 增长与 EventBridge、Converter 或权限错误关联分析。
导入 Grafana 大盘
下载Kafka Connect Grafana 大盘,在 Grafana 中导入 JSON,并选择已采集 Kafka Connect 指标的数据源;导入后确认数据源名称和标签与当前 Connect Cluster 的指标采集配置一致。限制条件
- EventBridge 单次
PutEvents请求最多包含 10 条事件,事件大小上限为 256 KiB;超过大小的单条事件不会通过 Connector 层重试机制解决,通常应使用 S3 外置或 DLQ 处理。 - Connector 对有效 Kafka 记录提供至少一次投递语义;EventBridge 成功后到 Kafka Connect 提交 Offset 前发生故障时,记录可能被再次发送,不能据此声明 exactly-once 或目标侧幂等。
- S3 外置只支持
$.detail.value下的 definite JSONPath,不支持 Key、Header 或 EventBridge 顶层字段,也不支持数组和通配符路径。 - 配置
aws.eventbridge.detail.types.mapper.class后,aws.eventbridge.detail.types不参与映射;两者不会合并生效。
常见问题
Task 启动时提示事件总线 ARN、区域或 Connector 标识无效
检查aws.eventbridge.connector.id 是否为空,aws.eventbridge.region 是否为 AWS SDK 支持的区域名称,以及 aws.eventbridge.eventbus.arn 是否为完整的事件总线 ARN。修正后重新验证 Connector 配置,并确认该 ARN 对应的区域与 aws.eventbridge.region 一致。
EventBridge 返回权限错误,记录没有送达
确认 Connector 实际使用的 AWS 身份来源。使用 IAM 角色时检查信任关系、角色 ARN 和 External ID;使用 Profile 时确认 Worker 能读取对应的 AWS 配置文件,并移除会覆盖 Profile 选择的 AWS 凭据环境变量。目标身份至少需要对目标事件总线拥有events:PutEvents 权限。
detail-type 没有使用消息中的事件类型
如果使用默认映射,detail-type 来自 Topic 或 aws.eventbridge.detail.types,不会自动读取 Value 中的字段。若要从 JSON Value 提取事件类型,将 aws.eventbridge.detail.types.mapper.class 设置为 software.amazon.event.kafkaconnector.mapping.JsonPathDetailTypeMapper,并配置一个 definite 的 aws.eventbridge.detail.types.jsonpathmapper.fieldref。提取结果必须是非空字符串,否则会回退到 Topic。
大消息仍然无法发送到 EventBridge
确认aws.eventbridge.offloading.default.s3.bucket 非空,aws.eventbridge.offloading.default.fieldref 以 $.detail.value 开头且确实匹配 Value 中的字段,并检查 AWS 身份是否具有 S3 PutObject 权限。若事件中没有匹配值或值为 null,Connector 会透传原事件,不会生成 S3 引用。
失败记录没有出现在 DLQ
确认同时配置了errors.tolerance=all 和 errors.deadletterqueue.topic.name,并检查 DLQ Topic 的创建、写入权限和副本因子。EventBridge 可重试错误会先按 aws.eventbridge.retries.max 重试;不是所有错误都会被归类为可报告的记录级失败,需结合 Task 日志判断是重试耗尽、转换失败还是任务级错误。