> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Server Sent Events Source Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Server Sent Events Source Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.github.cjmatta.kafka.connect.sse.ServerSentEventsSourceConnector
sse.uri=<sse-endpoint-url>
topic=<kafka-topic>
```

将 `<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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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 流的事件速率。
