概述
Redis Sink Connector 从 Kafka Topic 消费记录,并根据配置将记录写入 Redis。它可以把记录写入字符串、哈希、JSON 文档、Stream、List、Set、Sorted Set 或 RedisTimeSeries。redis.command 决定目标数据结构和记录键值的映射方式,redis.key 决定 Redis 键空间;集合类型使用 Topic 展开的键空间,字符串、哈希和 JSON 等非集合类型通常还会拼接 Kafka 记录键。
该 Sink 适用于把 Kafka 中的事件、状态或指标同步到 Redis,供低延迟查询、缓存、集合操作、流式消费或时间序列场景使用。投递语义为至少一次;选择会追加数据的 Redis 命令时,需要结合重放可能性设计目标键和业务去重策略。
前置条件
- 使用
JSONSET时,目标 Redis 部署需要提供 RedisJSON 能力;使用TSADD时,需要提供 RedisTimeSeries 能力。 - 使用 TLS、ACL 或双向 TLS 时,Connect Worker 必须能够读取配置的凭据、CA 证书、客户端证书和私钥文件。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka 和 Redis,并确认 Worker 可以访问目标 Redis、消费指定 Topic 以及使用目标 Redis 命令。具体准备和管理操作请参阅 管理 Connector。orders 替换为要消费的 Kafka Topic,将 <redis-host> 替换为 Redis 地址。此示例把字符串记录值写入 Redis String;非空记录键会生成 <topic>:<record-key> 形式的目标键。将配置提交到 Connect Cluster 后,检查 Connector 和 Task 状态,并确认目标 Redis 中出现预期键值。
配置
Redis 连接
redis.uri
用于创建 Redis 连接的 URI,支持 redis:// 和 rediss:// 方案。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:非空时优先于
redis.host和redis.port;为空时使用主机和端口创建连接。
redis.host
Redis 主机地址。
- 类型:
string - 默认值:
localhost - 重要级别:高
- 有效值 / 注意事项:仅在
redis.uri为空时使用;填写 Worker 可访问的主机名或地址。
redis.port
Redis TCP 端口。
- 类型:
int - 默认值:
6379 - 重要级别:高
- 有效值 / 注意事项:仅在
redis.uri为空时与redis.host一起使用;填写有效的 Redis TCP 端口。
redis.cluster
是否使用 Redis Cluster 客户端。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true使用 Cluster 客户端,false使用普通 Redis 客户端;使用redis.multiexec=true时不要选择 Cluster 客户端。
redis.timeout
Redis 命令超时时间,单位为秒。
- 类型:
long - 默认值:
60 - 重要级别:中
- 有效值 / 注意事项:用于 Redis 客户端和批量写入等待;根据 Redis 响应时间和业务延迟要求调整。
redis.pool
Redis Writer 使用的最大连接池大小。
- 类型:
int - 默认值:
8 - 重要级别:中
- 有效值 / 注意事项:连接池大小会影响并发写入能力;应结合 Task 数量、Redis 连接限制和 Worker 资源调整。
身份认证与 TLS
redis.username
用于 Redis ACL 认证的用户名。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:与非空的
redis.password一起使用。
redis.password
Redis 密码。
- 类型:
password - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:非空时应用到 Redis URI;不要将真实密码写入文档或日志。
redis.tls
是否启用 Redis TLS。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true启用 TLS;需要根据服务端证书配置可信 CA 或客户端证书。
redis.insecure
是否在 TLS 连接中关闭对端证书校验。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅在
redis.tls=true时生效;true会关闭对端校验,仅适用于已经评估过风险的环境。
redis.cacert
用于校验 Redis 服务端的 X.509 CA 证书文件。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:填写 Worker 可读的 CA 证书文件路径;为空时使用默认信任配置。
redis.key.file
用于双向 TLS 的 PKCS#8 私钥 PEM 文件。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:填写 Worker 可读的私钥文件路径;设置后还应设置对应的
redis.key.cert。
redis.key.cert
与 redis.key.file 配对的 X.509 客户端证书链 PEM 文件。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:在设置
redis.key.file时提供对应的证书链文件。
redis.key.password
redis.key.file 私钥的密码。
- 类型:
password - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:私钥已加密时填写私钥密码;未加密私钥使用空值。
目标数据映射
redis.command
选择写入 Redis 的数据结构或操作。
- 类型:
string - 默认值:
XADD - 重要级别:高
- 有效值 / 注意事项:可选
HSET、JSONSET、TSADD、SET、XADD、LPUSH、RPUSH、SADD、ZADD或DEL,大小写必须匹配。HSET和XADD需要 Struct 或 Map 值;ZADD的值为数字分数;TSADD的记录键为毫秒时间戳、记录值为数值样本。
redis.key
目标 Redis 键空间的格式字符串。
- 类型:
string - 默认值:
${topic} - 重要级别:中
- 有效值 / 注意事项:
${topic}会替换为来源 Kafka Topic。集合类型只使用展开后的键空间;其他类型通常再使用redis.separator拼接记录键。为空时,非集合类型可直接使用 Kafka 记录键作为目标键。
redis.separator
非集合目标中键空间与 Kafka 记录键之间的分隔符。
- 类型:
string - 默认值:
: - 重要级别:中
- 有效值 / 注意事项:应用于非集合目标;配置值会先去除首尾空白。
redis.charset
编码 Redis 键和值字符串时使用的字符集。
- 类型:
string - 默认值:
Charset.defaultCharset().name()返回的字符集 - 重要级别:高
- 有效值 / 注意事项:必须是 JVM 可识别的字符集名称;Kafka 记录键和值的转换方式还取决于 Converter 配置。
写入交付
redis.multiexec
是否使用 Redis MULTI/EXEC 事务执行批量写入。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅支持
XADD、LPUSH、RPUSH、SADD和ZADD;不支持与HSET、JSONSET、TSADD、SET或DEL组合使用。
redis.wait.replicas
写入后等待确认的 Redis 副本数量。
- 类型:
int - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:
0表示不等待副本;大于0时会执行 Redis WAIT,并受redis.wait.timeout限制。
redis.wait.timeout
Redis WAIT 命令的超时时间,单位为毫秒。
- 类型:
long - 默认值:
1000 - 重要级别:中
- 有效值 / 注意事项:与
redis.wait.replicas配合使用;只有设置等待副本数量时才会产生等待行为。
Kafka 消费与任务
topics
该 Sink 消费的 Kafka Topic 列表。
- 类型:
list - 默认值:
"" - 重要级别:高
- 有效值 / 注意事项:与
topics.regex互斥;二者必须且只能有一个非空。
topics.regex
用于匹配 Kafka Topic 的 Java 正则表达式。
- 类型:
string - 默认值:
"" - 重要级别:高
- 有效值 / 注意事项:与
topics互斥;表达式必须是有效的 Java 正则,并且不能匹配 DLQ Topic。
tasks.max
允许启动的最大 Sink Task 数量。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1;实际并行度还受 Kafka 分区数和分配结果影响,增加 Task 不提供跨分区的全局顺序。
key.converter
将 Kafka 消息键反序列化为 Connect 数据的 Converter。
- 类型:
class - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:使用可实例化的 Kafka Connect Converter;字符串或字节键分别使用
org.apache.kafka.connect.storage.StringConverter或org.apache.kafka.connect.converters.ByteArrayConverter。
value.converter
将 Kafka 消息值反序列化为 Connect 数据的 Converter。
- 类型:
class - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:使用可实例化的 Kafka Connect Converter;目标命令要求的值类型必须与 Converter 和 SMT 输出一致。
错误处理
errors.tolerance
控制 Kafka Connect 是否容忍记录处理错误。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可选
none或all;它是 Kafka Connect 的错误处理设置,不会改变 Redis 命令的数据映射要求。
errors.retry.timeout
Kafka Connect 错误重试的总时长,单位为毫秒。
- 类型:
long - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:
0表示不进行 Connect 错误重试,-1表示无限重试,其他值表示重试时长;它与redis.timeout的 Redis 命令超时不同。
errors.retry.delay.max.ms
Kafka Connect 错误重试之间的最大延迟,单位为毫秒。
- 类型:
long - 默认值:
60000 - 重要级别:中
- 有效值 / 注意事项:仅在启用 Connect 错误重试时生效,实际延迟可能包含抖动。
errors.deadletterqueue.topic.name
用于保存错误记录的 DLQ Topic 名称。
- 类型:
string - 默认值:
"" - 重要级别:中
- 有效值 / 注意事项:为空时禁用 DLQ;配置的 DLQ Topic 不得被
topics消费或被topics.regex匹配。
最佳实践
首次接入时先写入 Redis String
适用业务场景:需要先把 Kafka 中的字符串或字节值写入 Redis,验证 Topic、连接和键空间,再逐步扩展到其他 Redis 数据结构。SET 时,非空记录键会参与生成目标键;同一目标键的后续写入会覆盖已有值。该配置与快速开始相同,适合作为首次接入的最小验证路径。
将结构化记录写入 Hash
适用业务场景:Kafka 值已经由 Avro 或 JSON Converter 转换为 Struct 或 Map,需要把字段写入 Redis Hash,便于按字段读取。HSET 需要 Struct 或 Map 值,字段名作为 Hash 字段;非空 redis.key 会与记录键组合生成非集合目标键。使用 JSON Converter 时,确保输入记录的 JSON 结构能转换为 Map;需要 schemaful 数据时,应改用相应的结构化 Converter。
将事件追加到 Redis Stream
适用业务场景:需要把 Kafka 事件追加到 Redis Stream,供 Redis Stream 消费者继续处理,而不是只保留某个键的最新状态。XADD 使用 Struct 或 Map 的字段作为 Stream 消息体;redis.multiexec=true 只适用于支持 MULTI/EXEC 的命令集合,其中包括 XADD。该模式是追加写入,至少一次投递或重放可能造成重复 Stream 消息;Task 数量、Kafka 分区或重启也不构成跨分区的全局顺序保证,不能作为精确一次保证。
监控
监控内容
关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;Redis 连接错误、命令超时和批量写入延迟应结合 Task 错误与重试观察;仅在启用了相应错误处理时关注 DLQ 活动。导入 Grafana 大盘
确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。限制条件
redis.multiexec=true只支持XADD、LPUSH、RPUSH、SADD和ZADD,且不能与 Redis Cluster 客户端组合使用。- Connector 不提供通用的 Schema 到 Redis 数据结构映射;
HSET和XADD需要 Struct 或 Map,ZADD和TSADD需要数值输入。 - Connector 使用至少一次投递语义,不提供精确一次处理、Connector 级幂等或重复抑制;List 和 Stream 等追加型命令在重放时可能再次追加。
- Connector 没有自有的批量大小配置;每次 Kafka Connect 交付给 Task 的记录集合会作为一个 Writer 批次处理,过大的 Worker 批次可能受 Redis 超时和服务端容量影响。
topics与topics.regex互斥,且 Sink 必须配置其中一个非空;DLQ Topic 不能同时被该 Sink 消费。
常见问题
Task 启动后无法连接 Redis
检查redis.uri 是否为空以及 redis.host、redis.port 是否正确;非空 URI 会优先于主机和端口。启用 TLS 时检查 redis.tls、CA 文件、客户端证书和私钥路径是否对 Worker 可读;启用 ACL 时检查 redis.username 与 redis.password 是否匹配。修正连接配置后,确认 Task 恢复运行并观察 Redis 端是否收到连接。
写入 Hash 或 Stream 时出现值类型错误
HSET 和 XADD 需要 Struct 或 Map 值。检查 value.converter 及 SMT 输出,确认记录值不是普通字符串或不支持的对象;JSON Converter 应输出可转换为 Map 的结构。若业务数据本身是字符串,应改用 SET 或 JSONSET,不要依赖 Connector 自动把字符串解析为 Hash 字段。
记录写入的 Redis 键与预期不符
检查redis.key、redis.separator 和 Kafka 记录键。${topic} 会替换为来源 Topic;集合类型只使用展开后的 redis.key,而非集合类型通常还会追加记录键。记录键必须能转换为字符串或字节;空或不支持的记录键可能导致写入前转换失败。
重启后 Redis 中出现重复事件
该 Sink 使用至少一次投递,目标写入完成但 Offset 尚未持久化时发生故障,重启后可能重放记录。SET、HSET、JSONSET 和 DEL 通常对同一目标键产生覆盖或删除效果,而 LPUSH、RPUSH 和 XADD 可能再次追加;需要避免重复时,应在业务层设计幂等键或去重逻辑。
开启 redis.multiexec 后配置无效
确认 redis.command 是 XADD、LPUSH、RPUSH、SADD 或 ZADD,并确认没有设置 redis.cluster=true。其他命令会被配置校验拒绝;Redis Cluster 客户端也不支持该事务路径。