Skip to main content

概述

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.hostredis.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
  • 重要级别:高
  • 有效值 / 注意事项:可选 HSETJSONSETTSADDSETXADDLPUSHRPUSHSADDZADDDEL,大小写必须匹配。HSETXADD 需要 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
  • 重要级别:中
  • 有效值 / 注意事项:仅支持 XADDLPUSHRPUSHSADDZADD;不支持与 HSETJSONSETTSADDSETDEL 组合使用。

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.StringConverterorg.apache.kafka.connect.converters.ByteArrayConverter

value.converter

将 Kafka 消息值反序列化为 Connect 数据的 Converter。
  • 类型class
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:使用可实例化的 Kafka Connect Converter;目标命令要求的值类型必须与 Converter 和 SMT 输出一致。

错误处理

errors.tolerance

控制 Kafka Connect 是否容忍记录处理错误。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 noneall;它是 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 只支持 XADDLPUSHRPUSHSADDZADD,且不能与 Redis Cluster 客户端组合使用。
  • Connector 不提供通用的 Schema 到 Redis 数据结构映射;HSETXADD 需要 Struct 或 Map,ZADDTSADD 需要数值输入。
  • Connector 使用至少一次投递语义,不提供精确一次处理、Connector 级幂等或重复抑制;List 和 Stream 等追加型命令在重放时可能再次追加。
  • Connector 没有自有的批量大小配置;每次 Kafka Connect 交付给 Task 的记录集合会作为一个 Writer 批次处理,过大的 Worker 批次可能受 Redis 超时和服务端容量影响。
  • topicstopics.regex 互斥,且 Sink 必须配置其中一个非空;DLQ Topic 不能同时被该 Sink 消费。

常见问题

Task 启动后无法连接 Redis

检查 redis.uri 是否为空以及 redis.hostredis.port 是否正确;非空 URI 会优先于主机和端口。启用 TLS 时检查 redis.tls、CA 文件、客户端证书和私钥路径是否对 Worker 可读;启用 ACL 时检查 redis.usernameredis.password 是否匹配。修正连接配置后,确认 Task 恢复运行并观察 Redis 端是否收到连接。

写入 Hash 或 Stream 时出现值类型错误

HSETXADD 需要 Struct 或 Map 值。检查 value.converter 及 SMT 输出,确认记录值不是普通字符串或不支持的对象;JSON Converter 应输出可转换为 Map 的结构。若业务数据本身是字符串,应改用 SETJSONSET,不要依赖 Connector 自动把字符串解析为 Hash 字段。

记录写入的 Redis 键与预期不符

检查 redis.keyredis.separator 和 Kafka 记录键。${topic} 会替换为来源 Topic;集合类型只使用展开后的 redis.key,而非集合类型通常还会追加记录键。记录键必须能转换为字符串或字节;空或不支持的记录键可能导致写入前转换失败。

重启后 Redis 中出现重复事件

该 Sink 使用至少一次投递,目标写入完成但 Offset 尚未持久化时发生故障,重启后可能重放记录。SETHSETJSONSETDEL 通常对同一目标键产生覆盖或删除效果,而 LPUSHRPUSHXADD 可能再次追加;需要避免重复时,应在业务层设计幂等键或去重逻辑。

开启 redis.multiexec 后配置无效

确认 redis.commandXADDLPUSHRPUSHSADDZADD,并确认没有设置 redis.cluster=true。其他命令会被配置校验拒绝;Redis Cluster 客户端也不支持该事务路径。