概述
Server Sent Events Source Connector 连接一个支持 Server-Sent Events(SSE)的 HTTP 端点,持续接收事件并写入指定的 Kafka Topic。它适合将服务状态更新、通知流或其他由 SSE 长连接发布的实时数据接入 Kafka,供下游流处理、存储或分发系统使用。 每个 SSE 事件对应一条 Kafka 记录。记录值使用固定结构,包含event、id 和 data 三个字符串字段;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 - 重要级别:中
- 有效值 / 注意事项:可设为
true或false。启用后,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 - 重要级别:低
- 有效值 / 注意事项:可设为
true或false。启用时,Connector 会设置 Jersey 的 gzip 编码属性,但显式构造的Accept-Encoding请求头没有接入实际 SSE 握手路径。在目标 SSE 服务上验证压缩协商和响应解压行为。
连接节流
rate.limit.requests.per.second
控制 Connector 打开或重新打开 SSE 连接前的最小请求间隔。
- 类型:
double - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:正数用于计算连接打开请求之间的等待时间;
null、0或负数不启用该等待。它不限制 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内容解析;SSEdata始终作为字符串字段写入固定结构。 - 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 流的事件速率。