使用KSqlDB/Kafka Streams按优先级拆分源Topic消息时,如何保持原消息格式不变?
KSqlDB/Kafka Streams按优先级拆分源Topic消息时,如何保持原消息格式不变?
嗨,我来帮你搞定这个问题!你现在遇到的核心问题是拆分后data字段被转成了转义的JSON字符串,这是因为你在定义源流的时候把data指定成了varchar类型——KSqlDB会把嵌套的JSON对象当成普通字符串处理,输出时自然就会转义啦。
下面给你两种靠谱的解决方案,都能让原消息完整无损地分流:
方案一:让KSqlDB自动处理整个JSON消息(最简单)
不需要手动拆分字段,直接让KSqlDB把整个消息当作JSON对象处理,这样就能直接基于data.priority过滤,同时保留原消息结构:
- 创建源流(不指定具体字段,让KSqlDB自动识别JSON结构):
CREATE STREAM source_stream WITH ( KAFKA_TOPIC = 'topic_source', VALUE_FORMAT = 'JSON' );
- 创建优先级1的分流流:
CREATE STREAM source_stream_priority_1 WITH ( KAFKA_TOPIC = 'topic_source_priority_1', VALUE_FORMAT = 'JSON' ) AS SELECT * FROM source_stream WHERE EXTRACTJSONFIELD(VALUE, '$.data.priority') = 1;
- 创建优先级2的分流流:
CREATE STREAM source_stream_priority_2 WITH ( KAFKA_TOPIC = 'topic_source_priority_2', VALUE_FORMAT = 'JSON' ) AS SELECT * FROM source_stream WHERE EXTRACTJSONFIELD(VALUE, '$.data.priority') = 2;
注意:这里要把过滤条件写成数字
2,而不是字符串'2',因为原消息里的priority是数字类型,字符串匹配会失效哦。
方案二:手动指定字段但用JSON类型存储data
如果你需要明确指定顶层字段(比如id、time),可以把data定义为KSqlDB的JSON类型,这样它会保留嵌套结构,不会被转成字符串:
- 创建源流:
CREATE STREAM source_stream ( id VARCHAR, time VARCHAR, data JSON ) WITH ( KAFKA_TOPIC = 'topic_source', VALUE_FORMAT = 'JSON' );
- 分流时可以用更简洁的语法提取
priority:
CREATE STREAM source_stream_priority_1 WITH ( KAFKA_TOPIC = 'topic_source_priority_1', VALUE_FORMAT = 'JSON' ) AS SELECT * FROM source_stream WHERE data->>'priority' = 1; CREATE STREAM source_stream_priority_2 WITH ( KAFKA_TOPIC = 'topic_source_priority_2', VALUE_FORMAT = 'JSON' ) AS SELECT * FROM source_stream WHERE data->>'priority' = 2;
用这两种方法,目标Topic里的消息就会和源Topic完全一致,data字段依然是嵌套的JSON结构,不会出现转义字符串的问题啦!
备注:内容来源于stack exchange,提问作者DmitrySpb
相关产品推荐
相关产品推荐

