如何为ClickHouse Kafka Engine导出的Kafka消息设置Key?
解决方案
要给Kafka消息设置Key,你需要利用ClickHouse Kafka引擎的kafka_message_key_column配置项,具体操作如下:
- 修改物化视图定义,在Kafka引擎的SETTINGS中添加
kafka_message_key_column = '_key',指定用_key列作为Kafka消息的Key;同时调整SELECT语句,确保_key列被单独指定,其余列组成消息Value:
CREATE MATERIALIZED VIEW myView on cluster '{cluster}' ENGINE = Kafka('kafka:9092', 'MyTopic', 'MyConsumerGroup', 'JSONEachRow') SETTINGS kafka_thread_per_consumer = 0, kafka_num_consumers = 2, kafka_message_key_column = '_key' -- 指定作为Kafka消息Key的列 POPULATE AS select toYYYYMMDD(source_day) as _key, -- 该列值将作为Kafka消息Key lastUpdateTime, groupId, source_day, sourceId, languages from MyTable FORMAT JSONEachRow;
核心原理:
kafka_message_key_column参数会让ClickHouse将指定列的值提取为Kafka消息的Key,且不会把该列包含在消息的Value中;- 修改后,Kafka收到的消息会自动将
_key的值设为Key,Value由剩余字段组成,完全符合你的预期效果。
注意事项:
- 若已创建过同名物化视图,需先删除原视图再重新创建;
_key列的数据类型需为Kafka支持的类型(如字符串、数值类型均可,ClickHouse会自动转换为字节流作为Key)。
内容的提问来源于stack exchange,提问作者Sasha Sasha
相关产品推荐
相关产品推荐

