Skip to main content

概述

Server Sent Events Source Connector 连接一个支持 Server-Sent Events(SSE)的 HTTP 端点,持续接收事件并写入指定的 Kafka Topic。它适合将服务状态更新、通知流或其他由 SSE 长连接发布的实时数据接入 Kafka,供下游流处理、存储或分发系统使用。 每个 SSE 事件对应一条 Kafka 记录。记录值使用固定结构,包含 eventiddata 三个字符串字段;data 保留为字符串,不会自动解析其中的 JSON 或 XML。记录键为空,SSE 事件 ID 也不会作为 Kafka Connect Offset 保存。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和无需认证的 SSE 端点,并确认网络连通和访问权限。具体准备和管理操作请参阅 管理 Connector
<sse-endpoint-url> 替换为 Connect Worker 可访问的 SSE 端点,将 <kafka-topic> 替换为接收事件的 Kafka Topic。应用配置后,应同时检查 Connector 和 Task 状态,并从目标 Topic 验证实际事件;Task 处于运行状态本身不代表 SSE 连接已经产生数据。

配置

Kafka Connect 框架

connector.class

要加载的 Server Sent Events Source Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 com.github.cjmatta.kafka.connect.sse.ServerSentEventsSourceConnector
  • 必填:是

tasks.max

Kafka Connect 请求 Connector 创建的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须至少为 1。此 Connector 始终只创建一个 Task;增大该值不会增加 SSE 连接数或消费并发。

SSE 连接与输出

sse.uri

Task 要连接的 SSE 流地址。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用每个候选 Worker 都可访问的 SSE 端点。Connector 不在配置解析阶段校验 URI 格式、协议或可达性。
  • 必填:是

topic

接收所有 SSE 事件的 Kafka Topic。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:一个 Connector 实例只能写入一个固定 Topic,不支持按事件内容配置 Topic 列表或路由规则。
  • 必填:是

HTTP 身份验证与请求头

http.basic.auth

控制是否读取 HTTP Basic Authentication 的用户名和密码。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:可设为 truefalse。启用后,Task 会以字符串方式读取 password 类型的配置值,可能在启动时发生类型转换错误;即使成功构造 Authorization 请求头,当前实现也没有把对应的请求 builder 交给实际 SSE 握手。不要将其作为可用的认证方案。

http.basic.auth.username

HTTP Basic Authentication 的用户名。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 http.basic.auth=true 时读取;Connector 不校验非空值,也不校验它与密码是否同时配置。受 Basic Auth 实现限制,不应依赖此配置接入需要认证的端点。

http.basic.auth.password

HTTP Basic Authentication 的密码。
  • 类型password
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅在 http.basic.auth=true 时读取。启用 Basic Auth 并提供密码可能导致 Task 启动时发生类型转换错误;不要在明文配置、日志或文档中暴露真实密码。

http.header.*

读取自定义 HTTP 请求头配置;http.header. 后的名称作为请求头名称。
  • 类型dynamic-prefix string values
  • 默认值{}
  • 重要级别:中
  • 有效值 / 注意事项:例如 http.header.User-Agent 表示 User-Agent 请求头。Connector 不校验请求头名称和值;值可能包含敏感凭证。当前实现把这些请求头添加到一个未交给 SseEventSource 的请求 builder,因此不要依赖此配置完成认证、访问控制或 User-Agent 覆盖。

HTTP 客户端

compression.enabled

控制 Jersey 客户端的 gzip 编码处理属性。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:可设为 truefalse。启用时,Connector 会设置 Jersey 的 gzip 编码属性,但显式构造的 Accept-Encoding 请求头没有接入实际 SSE 握手路径。在目标 SSE 服务上验证压缩协商和响应解压行为。

连接节流

rate.limit.requests.per.second

控制 Connector 打开或重新打开 SSE 连接前的最小请求间隔。
  • 类型double
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:正数用于计算连接打开请求之间的等待时间;null0 或负数不启用该等待。它不限制 SSE 事件吞吐,也不是 Worker 全局请求速率限制。

rate.limit.max.concurrent

声明最大并发连接数。
  • 类型int
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:当前实现虽然接受并保存该值,但不会据此限制连接、Task 或请求并发,不应将其作为并发控制手段。

重连参数

retry.backoff.initial.ms

设置空闲健康检查触发重连时的初始等待时间,并参与 SSE 客户端自身重连间隔的计算。
  • 类型long
  • 默认值2000
  • 重要级别:低
  • 有效值 / 注意事项:单位为毫秒。Connector 不校验非负值;SSE 客户端自身使用的间隔不会超过 2000 毫秒,而空闲健康检查重连使用配置值作为第一次指数退避等待。

retry.backoff.max.ms

设置空闲健康检查触发重连时指数退避的最大等待时间。
  • 类型long
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:单位为毫秒。Connector 不校验非负值,也不要求该值大于或等于 retry.backoff.initial.ms;它不控制 SSE 客户端自身的全部重连路径。

retry.max.attempts

限制 Connector 管理的空闲健康检查重连次数。
  • 类型int
  • 默认值-1
  • 重要级别:低
  • 有效值 / 注意事项-1 表示不限次数,0 表示拒绝第一次此类重连;小于 -1 的值不会被视为不限次数。该配置不保证覆盖异步 SSE 错误、所有 HTTP 状态或 Task 失败后的恢复。

监控

监控内容

关注 Kafka Connect 通用健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;同时观察 SSE 连接错误及队列增长相关日志,避免源端持续快于 Kafka 导致堆内存压力;仅在启用了相应错误处理时关注 DLQ 活动,但不要依赖 DLQ 捕获 Connector 内部的连接和事件处理错误。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 一个 Connector 实例只连接一个 SSE URI、写入一个固定 Topic,并且始终只创建一个 Task;如需连接多个 SSE 端点,需要分别创建 Connector 实例。
  • Connector 不提供事件类型过滤或 data 内容解析;SSE data 始终作为字符串字段写入固定结构。
  • Connector 不把 SSE 事件 ID 保存为 Kafka Connect Offset,也不提供持久化去重;Task 或 Worker 重启及重连可能产生重复或遗漏,不能声明 exactly-once、无重复、无遗漏或无条件 at-least-once。
  • SSE 事件先进入无容量上限的内存队列,单次 poll 也没有记录数或字节数上限;源端持续快于 Kafka 时可能产生堆内存压力。
  • Basic Auth 存在密码类型转换失败路径;Authorization、自定义请求头和显式 Accept-Encoding 被添加到未交给 SseEventSource 的请求 builder,不应依赖这些配置访问受保护端点或覆盖 User-Agent。
  • rate.limit.max.concurrent 不执行并发限制;rate.limit.requests.per.second 只影响 Connector 主动打开连接前的等待,不限制事件吞吐或 SSE 客户端内部的全部重连。
  • 重连参数仅覆盖特定连接生命周期路径;它们不能保证所有初始连接失败、异步错误或 HTTP 响应都会自动重试并恢复。
  • 缺少 event 名称的 SSE 消息可能因其在一次批次中的位置不同而被跳过或以 event=unknown 输出。

常见问题

Connector 和 Task 显示运行中,但目标 Topic 没有消息

Task 启动成功不一定表示 SSE 流已连接并持续产生事件。确认 sse.uri 可从运行 Task 的 Worker 访问,端点返回有效的 SSE 响应且正在发布带数据的事件;检查 Worker 日志中的连接、协议和异步错误,并从目标 Topic 验证输出。对于需要 Basic Auth 或自定义请求头的端点,不要假设当前配置能够完成握手,优先使用无需认证的受控端点定位问题。

重启后出现重复消息或缺少一段事件

Connector 不保存可恢复的 SSE 消费位置,也不把事件 ID 写入 Kafka Connect Offset。即使 SSE 服务支持 Last-Event-ID,Task 或 Worker 重启后 Connector 也没有可用于恢复的持久化事件 ID。在下游使用稳定业务标识进行幂等处理或去重,并监控重启窗口内的重复与数据缺口;不要依赖 Offset 提交消除此风险。

增大 tasks.max 后仍然只有一个 Task

这是该 Connector 的任务分配边界。它始终返回一个 Task 配置,因此提高 tasks.max 不会增加 SSE 连接数或吞吐。需要接入多个端点时,为每个端点创建独立 Connector 实例;单个端点的容量规划应结合源端事件速率、Kafka 写入能力和 Worker 内存进行验证。

Worker 内存持续增长

SSE 回调会把事件写入无界内存队列,而 Connector 不提供背压或批次上限。检查源端事件速率、Kafka 生产延迟、Task 错误和队列相关日志;降低源端发送速率、恢复 Kafka 写入能力或拆分不同端点到独立 Connector 实例,并为 Worker 设置内存与告警阈值。rate.limit.requests.per.second 只限制打开连接的频率,不能降低已建立 SSE 流的事件速率。