> ## 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.

# Redis Stream Source Connector

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

## 概述

Redis Stream Source Connector 从一个 Redis Stream 读取消息，并将每条 Stream 消息转换为 Kafka Connect SourceRecord 后写入 Kafka Topic。它位于 Redis 与 Kafka 之间，适合把 Redis Stream 中的业务事件、任务状态或实时通知接入 Kafka 数据链路。

Connector 使用 Redis 消费者组分发消息。Kafka 消息 Key 是 Redis Stream 消息 ID，Value 是包含 `id`、`stream` 和字符串字段映射 `body` 的结构；目标 Topic 可以固定指定，也可以使用 `${stream}` 按 Redis Stream 名称生成。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和 Redis 服务，确认 Connect Worker 可以访问 Redis 和 Kafka。若 Kafka 集群禁用了 Topic 自动创建，请先创建用于接收记录的目标 Topic。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=com.redis.kafka.connect.RedisStreamSourceConnector
redis.host=<redis-host>
redis.stream.name=<redis-stream-name>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

将 `<redis-host>` 替换为 Redis 地址，将 `<redis-stream-name>` 替换为要读取的 Redis Stream Key。示例使用默认 Redis 端口 `6379`、默认消费者组和 `at-least-once` 交付模式；目标 Topic 默认与 Redis Stream 同名。需要认证、TLS 或固定目标 Topic 时，添加对应配置后再创建 Connector。

## 配置

### Redis 连接

#### `redis.cluster`

选择 Redis 单节点客户端或 Redis Cluster 客户端。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：`false` 使用单节点客户端，`true` 使用 Cluster 客户端；地址或 URI 必须与所选部署模式匹配。

#### `redis.host`

Redis 主机名或地址。

* **类型**：`string`
* **默认值**：`localhost`
* **重要级别**：高
* **有效值 / 注意事项**：仅在 `redis.uri` 为空时与 `redis.port` 一起使用；该地址必须能从 Connect Worker 解析和访问。

#### `redis.port`

Redis TCP 端口。

* **类型**：`int`
* **默认值**：`6379`
* **重要级别**：高
* **有效值 / 注意事项**：仅在 `redis.uri` 为空时与 `redis.host` 一起使用。ConfigDef 不校验端口范围，应提供 Redis 实际监听的有效端口。

#### `redis.uri`

使用 Redis URI 指定连接端点。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空值优先于 `redis.host` 和 `redis.port`，支持客户端接受的 `redis://` 和 `rediss://` URI。URI 可能包含凭据，应按敏感信息管理；非空的 `redis.password` 会在端点选择后应用单独配置的凭据。

#### `redis.timeout`

Redis 命令超时时间，单位为秒。

* **类型**：`long`
* **默认值**：`60`
* **重要级别**：中
* **有效值 / 注意事项**：使用适合 Redis 响应延迟的正数。ConfigDef 不校验范围，无效时长可能在客户端构建或运行时失败。

### Redis 认证

#### `redis.username`

Redis ACL 用户名。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `redis.password` 非空时使用。留空表示使用仅密码认证。

#### `redis.password`

Redis 密码。

* **类型**：`password`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空值启用独立凭据配置；同时设置 `redis.username` 时使用 ACL 用户名和密码认证，否则使用仅密码认证。请通过受控的 Secret 或 Config Provider 提供该值。

### TLS

#### `redis.tls`

为使用 `redis.host` 和 `redis.port` 构造的 Redis 连接显式启用 TLS。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：设置为 `true` 时启用 TLS。`redis.uri` 使用 `rediss://` 时，即使此项保持 `false` 也会建立 TLS 连接；但 `redis.insecure` 只有在此配置显式为 `true` 时才会被读取。

#### `redis.insecure`

关闭 TLS 对端证书校验。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `redis.tls=true` 时生效。生产环境不要启用；只配置 `rediss://` 而未设置 `redis.tls=true` 时，该值不会影响证书校验分支。

#### `redis.cacert`

指定用于校验 Redis 服务端证书的 CA 证书文件。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空值必须指向 Connect Worker 可读的 X.509 CA 证书文件。

#### `redis.key.file`

指定 TLS 客户端私钥文件。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：非空值启用客户端 Key Manager，文件必须是 Connect Worker 可读的 PKCS#8 PEM 私钥；同时需要提供匹配的 `redis.key.cert`。

#### `redis.key.cert`

指定 TLS 客户端证书链文件。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `redis.key.file` 非空时读取，必须指向与私钥匹配且 Connect Worker 可读的 PEM X.509 证书链。

#### `redis.key.password`

指定 TLS 客户端私钥密码。

* **类型**：`password`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在 `redis.key.file` 非空时读取。未加密私钥保持为空；加密私钥应通过受控的 Secret 或 Config Provider 提供密码。

### Stream 选择与路由

#### `redis.stream.name`

指定要读取的 Redis Stream Key。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：一个 Connector 实例只读取一个 Stream。必须提供可用的 Redis Key；显式空字符串虽然能通过 ConfigDef 字符串解析，但不能形成有效部署。
* **必填**：是

#### `topic`

指定 SourceRecord 写入的 Kafka Topic。

* **类型**：`string`
* **默认值**：`${stream}`
* **重要级别**：中
* **有效值 / 注意事项**：`${stream}` 会替换为 Redis Stream 名称；也可以填写固定 Topic 名称。最终名称必须符合 Kafka Topic 命名、存在或自动创建以及授权要求，其他占位符不会被解析。

#### `redis.stream.offset`

指定没有 Kafka Connect 已存 Source Offset 时使用的初始 Redis Stream Offset。

* **类型**：`string`
* **默认值**：`0-0`
* **重要级别**：中
* **有效值 / 注意事项**：应使用 Redis Stream Reader 接受的消息 ID。重启时已存的 Kafka Connect Offset 优先于该值；已有消费者组也不会因为修改此值而重置服务端游标。

### 消费者组与交付

#### `redis.stream.consumer.group`

指定 Redis Stream 消费者组名称。

* **类型**：`string`
* **默认值**：`kafka-consumer-group`
* **重要级别**：中
* **有效值 / 注意事项**：同一 Connector 的所有 Task 使用该消费者组。修改组名会切换 Connector 使用的 Redis Pending Entry 和确认状态。

#### `redis.stream.consumer.name`

指定 Redis Stream 消费者名称模板。

* **类型**：`string`
* **默认值**：`consumer-${task}`
* **重要级别**：中
* **有效值 / 注意事项**：`${task}` 会替换为数字 Task ID。`tasks.max` 大于 `1` 时应保留该占位符，避免多个 Task 使用相同消费者身份。

#### `redis.stream.delivery`

选择 Redis Stream 消息确认模式。

* **类型**：`string`
* **默认值**：`at-least-once`
* **重要级别**：中
* **有效值 / 注意事项**：仅接受区分大小写的 `at-least-once` 或 `at-most-once`。前者在 SourceTask Commit 时确认消息，可能重放；后者读取后立即确认，故障窗口可能丢失消息。两者都不提供跨 Redis 与 Kafka 的 exactly-once 事务。

### 读取节奏

#### `batch.size`

设置每次 Redis 读取请求的最大记录数。

* **类型**：`int`
* **默认值**：`500`
* **重要级别**：低
* **有效值 / 注意事项**：应使用正数。它限制单次读取请求的数量，不是跨多次 Poll 的 Pending Entry 或内存总上限；ConfigDef 不校验上下界。

#### `redis.stream.block`

设置 Redis `XREADGROUP` 的阻塞时长，单位为毫秒。

* **类型**：`long`
* **默认值**：`100`
* **重要级别**：低
* **有效值 / 注意事项**：使用 Redis 客户端支持的非负值。该值直接影响没有新消息时单次 Poll 的等待时间，ConfigDef 不校验范围。

### 兼容配置

#### `redis.pool`

继承的 Redis 连接池大小配置。

* **类型**：`int`
* **默认值**：`8`
* **重要级别**：中
* **有效值 / 注意事项**：该配置公开可设置，但 RedisStreamSourceConnector 0.9.1 的 Stream 读取路径没有使用它；调整该值没有确认的 Stream Source 调优效果。

### Kafka Connect 框架

#### `connector.class`

指定要加载的 Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必须设置为 `com.redis.kafka.connect.RedisStreamSourceConnector`。
* **必填**：是

#### `tasks.max`

设置 Kafka Connect 为 Connector 创建的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须至少为 `1`。每个 Task 都读取同一个 Redis Stream 和消费者组；多 Task 时应在 `redis.stream.consumer.name` 中保留 `${task}`。

#### `tasks.max.enforce`

控制框架是否强制执行 `tasks.max` 上限。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：Kafka Connect 已将该配置标记为弃用并计划在未来主要版本移除。保持默认值，并通过 `tasks.max` 设置规模。
* **已弃用**：是
* **替代项**：无

#### `key.converter`

指定 Redis Stream 消息 ID Key 的 Kafka Connect Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的 Key Converter。Connector 发出带 STRING Schema 的非空 Key，应选择兼容该 Schema 的 Converter。

#### `value.converter`

指定结构化消息 Value 的 Kafka Connect Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的 Value Converter。Connector 发出包含 `id`、`stream` 和 `MAP<STRING,STRING>` 类型 `body` 的 STRUCT，应选择支持这些 Schema 的 Converter。

## 最佳实践

### 生产接入时启用 TLS 和 ACL 认证

**适用业务场景**：Redis 位于生产网络或使用 ACL 管理访问，需要加密 Connect Worker 与 Redis 之间的连接，同时只授予 Connector 读取目标 Stream 和使用消费者组所需的权限。

**配置示例**：在快速开始配置基础上添加以下连接配置。

```properties theme={null}
connector.class=com.redis.kafka.connect.RedisStreamSourceConnector
redis.host=<redis-host>
redis.port=<redis-tls-port>
redis.tls=true
redis.cacert=<redis-ca-certificate-path>
redis.username=<redis-username>
redis.password=<redis-password>
redis.stream.name=<redis-stream-name>
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

**关键说明**：CA 文件必须在每个可能运行 Task 的 Worker 上可读，用户名和密码应通过受控的 Secret 或 Config Provider 注入。不要在生产环境设置 `redis.insecure=true`；如果 Redis 要求双向 TLS，再同时配置匹配的客户端私钥和证书链。

### 根据吞吐和空闲延迟调整读取节奏

**适用业务场景**：Connector 已稳定读取单个 Stream，需要在单次处理开销、吞吐和无新消息时的返回延迟之间做权衡。先从默认值观察吞吐、延迟和 Worker 资源，再小步调整。

**配置示例**：以下示例提高单次请求数量，并允许空闲读取等待更长时间。

```properties theme={null}
connector.class=com.redis.kafka.connect.RedisStreamSourceConnector
redis.host=<redis-host>
redis.stream.name=<redis-stream-name>
batch.size=1000
redis.stream.block=500
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

**关键说明**：`batch.size` 是 Redis 单次读取请求的数量上限，不是内存或 Pending Entry 总上限；`redis.stream.block` 会直接增加空闲 Poll 的等待时间。示例值不是通用最优值，应结合消息大小、目标 Topic 写入能力、Task 延迟和 Worker 内存逐步调整。

### 单个 Stream 出现持续积压时增加 Task

**适用业务场景**：单 Task 已成为读取瓶颈，Redis Stream 消费者组持续积压，而且 Kafka 和 Connect Worker 仍有处理余量。增加 Task 可以让多个 Redis 消费者共享同一消费者组处理新消息。

**配置示例**：以下示例为同一个 Stream 创建两个 Task，并保留按 Task 生成的消费者名称。

```properties theme={null}
connector.class=com.redis.kafka.connect.RedisStreamSourceConnector
redis.host=<redis-host>
redis.stream.name=<redis-stream-name>
tasks.max=2
redis.stream.consumer.name=consumer-${task}
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

**关键说明**：多个 Task 不会把不同 Stream 分配给不同 Task，它们仍读取同一个 Stream 和消费者组。0.9.1 中所有 Task 还共用同一个 Kafka Connect Source Offset Key，而 Redis Pending Entry 按消费者身份分别归属。扩展前后应保持消费者名称稳定且唯一，并在生产使用前验证重启、重平衡和缩容恢复；稳态分流结果不能证明每个 Task 都有独立的恢复检查点。多 Task 与多分区 Kafka Topic 也不提供跨 Task 或跨分区的全局顺序。

## 监控

### 监控内容

关注 Kafka Connect 通用健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号；同时观察 Redis 消费者组积压与 Pending Entry 变化。仅在部署启用了相应错误处理时关注 DLQ 活动。

### 导入 Grafana 大盘

确认 Kafka Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* 一个 Connector 实例只能读取一个 Redis Stream Key，不支持 Stream 列表或模式匹配。
* 已存在的 Redis 消费者组不会因修改 `redis.stream.offset` 而重置服务端游标；重启时 Kafka Connect 已存 Source Offset 也优先于该配置。
* Connector 不提供跨 Redis 确认与 Kafka 写入的 exactly-once 事务或重复抑制；`at-least-once` 模式可能重放消息，`at-most-once` 模式可能在确认后、写入 Kafka 前丢失消息。
* Connector 只恢复当前消费者身份拥有的 Pending Entry，不会自动通过 `XCLAIM` 或 `XAUTOCLAIM` 接管其他消费者名下的消息。
* 0.9.1 中所有 Task 使用同一个空 Source Partition，并共用一个 Kafka Connect Source Offset Key；不能据此假定每个 Redis 消费者拥有独立的恢复检查点，多 Task 的故障恢复需要单独验证。
* 0.9.1 在 `at-least-once` 模式成功 Commit 后不会清空 Task 内已收集的消息 ID；后续 Commit 会重复向 Redis 发送历史 `XACK` ID，且该内存列表会在 Task 生命周期内持续增长。长时间高吞吐运行可能增加 Task 内存和 Redis 确认开销，目前没有已验证的安全阈值。
* 顺序只限于单个 Task 的一次 Poll 结果；多个 Task 或多分区 Kafka Topic 不提供全局顺序。
* 消息 Value Schema 固定为字符串字段映射，不保留任意二进制字段值，也不会自动推断或演进字段类型。

## 常见问题

### 修改 `redis.stream.offset` 后，为什么没有从新位置开始读取？

Kafka Connect 已存的 Source Offset 会优先于 `redis.stream.offset`，而且 Connector 打开已存在的 Redis 消费者组时不会重置该组的服务端游标。先确认使用的消费者组、Kafka Connect Offset 和 Redis 组状态；需要建立独立读取基线时，应使用新的消费者组并谨慎处理原组的 Pending Entry，而不是只修改初始 Offset。

### 为什么 Kafka 中出现重复消息？

默认 `at-least-once` 模式在 Kafka Connect Commit 时确认 Redis 消息，Redis 读取、Kafka 写入和 Offset 提交之间没有跨系统事务。故障发生在写入后、确认或 Offset 持久化前时，消息可能在重启后重放。Kafka 消息 Key 是 Redis Stream 消息 ID，可由下游结合该 ID 实现幂等处理或去重。

### 缩减 Task 后，为什么旧消费者的 Pending Entry 没有被继续处理？

Connector 只读取当前消费者名称拥有的 Pending Entry，不会自动接管其他消费者名下的消息。检查 Redis 消费者组中的 Pending Entry 所有者；扩缩容时保持 `consumer-${task}` 生成的消费者名称稳定，并在移除消费者前处理其 Pending Entry，必要时通过 Redis 运维流程显式转移所有权。

### 为什么消息时间戳与 Redis Stream ID 中的时间不一致？

SourceRecord 时间戳取自 Connector 执行消息转换时的 Worker 时钟，不使用 Redis Stream ID 中编码的时间部分。需要事件时间时，应在 Stream 消息字段中明确写入业务时间，并由下游从 `body` 中提取和使用。
