使用KSQL的PARTITION BY后值字段丢失ID,如何保留ID并按其分区?
解决KSQL分区后Value丢失分区字段的问题
问题原因
当使用PARTITION BY ID时,KSQL默认会把ID字段作为消息的key(用于将消息路由到对应分区),并且默认不会将该字段保留在value中,这就是新生成Topic的Avro Schema里没有ID字段的原因。
解决方案
要同时实现按ID分区且保留ID在value中,你需要在SELECT语句里显式指定包含ID字段,而不是使用SELECT *。这样KSQL会同时把ID放在消息的key(用于分区)和value里。
修改后的KSQL语句如下:
CREATE STREAM TEST_STREAM_AVRO WITH (PARTITIONS=3, VALUE_FORMAT='AVRO') AS SELECT id, age, name FROM TEST_STREAM_JSON PARTITION BY ID;
效果验证
执行上述语句后,新Topic的Avro Schema会包含ID字段,示例如下:
{"fields": [ {"default": null, "name": "ID", "type": ["null", "int"] }, {"default": null, "name": "AGE", "type": ["null", "int"] }, {"default": null, "name": "NAME", "type": ["null", "string"] } ], "name": "KsqlDataSourceSchema", "namespace": "io.confluent.ksql.avro_schemas", "type": "record"}
内容的提问来源于stack exchange,提问作者user21169353
相关产品推荐
相关产品推荐

