概述
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。<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 和使用消费者组所需的权限。 配置示例:在快速开始配置基础上添加以下连接配置。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 生成的消费者名称。监控
监控内容
关注 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,不会自动通过
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 发送历史XACKID,且该内存列表会在 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 中提取和使用。