Skip to main content

概述

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 是包含 idstream 和字符串字段映射 body 的结构;目标 Topic 可以固定指定,也可以使用 ${stream} 按 Redis Stream 名称生成。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和 Redis 服务,确认 Connect Worker 可以访问 Redis 和 Kafka。若 Kafka 集群禁用了 Topic 自动创建,请先创建用于接收记录的目标 Topic。具体准备和管理操作请参阅 管理 Connector
<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.hostredis.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.hostredis.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-onceat-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 发出包含 idstreamMAP<STRING,STRING> 类型 body 的 STRUCT,应选择支持这些 Schema 的 Converter。

最佳实践

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

适用业务场景:Redis 位于生产网络或使用 ACL 管理访问,需要加密 Connect Worker 与 Redis 之间的连接,同时只授予 Connector 读取目标 Stream 和使用消费者组所需的权限。 配置示例:在快速开始配置基础上添加以下连接配置。
关键说明:CA 文件必须在每个可能运行 Task 的 Worker 上可读,用户名和密码应通过受控的 Secret 或 Config Provider 注入。不要在生产环境设置 redis.insecure=true;如果 Redis 要求双向 TLS,再同时配置匹配的客户端私钥和证书链。

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

适用业务场景:Connector 已稳定读取单个 Stream,需要在单次处理开销、吞吐和无新消息时的返回延迟之间做权衡。先从默认值观察吞吐、延迟和 Worker 资源,再小步调整。 配置示例:以下示例提高单次请求数量,并允许空闲读取等待更长时间。
关键说明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 生成的消费者名称。
关键说明:多个 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,在 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,不会自动通过 XCLAIMXAUTOCLAIM 接管其他消费者名下的消息。
  • 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 中提取和使用。