You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过RedisSinkConnector将Kafka消息按Key存入Redis列表?

问题解决方法

你的问题源于Redis Sink Connector默认行为与预期不符:它会将Kafka的topic名称作为Redis列表的key,把Kafka消息的key作为列表元素。需要添加两个关键配置参数修正这一行为:

修改后的完整Connector配置

curl --location 'http://localhost:8083/connectors/' \
--header 'Content-Type: application/json' \
--data '{
  "name": "redis-sink",
  "config": {
    "connector.class": "com.redis.kafka.connect.RedisSinkConnector",
    "tasks.max": "1",
    "topics": "t_location",
    "redis.type": "LIST",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "redis.key.expression": "${record.key}",
    "redis.value.expression": "${record.value}"
  }
}'

关键配置说明

  • redis.key.expression: 配置Redis列表的key来源为Kafka消息的key(如car-three),通过${record.key}表达式提取每条Kafka记录的key值。
  • redis.value.expression: 配置Redis列表的元素来源为Kafka消息的完整JSON value,通过${record.value}表达式提取每条Kafka记录的value内容。

修正后效果

配置生效后,Connector会对每条Kafka消息执行你期望的Redis命令:

LPUSH car-three {"timestamp":1683742700544,"distanceFromBase":104,"car_no":"car-three"}

内容的提问来源于stack exchange,提问作者Solmaz Oskouie

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 14:02:25