如何通过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
相关产品推荐
相关产品推荐

