Flink SQL Client读取Kafka消息Key时orderNumber字段为NULL问题
Flink SQL读取Kafka消息Key字段为NULL的问题
场景说明
我有一个名为Orders的Kafka主题,消息内容如下:
Key: {"orderNumber":"1234"} Value: {"orderDate":"20250528","productId":"Product123"}
使用以下Kafka控制台消费者命令可以正常获取消息的Key和Value:
kafka-console-consumer.bat --property "print.key=true" --topic Orders --from-beginning --bootstrap-server xxxxxx
尝试的Flink SQL表定义
我尝试创建如下Flink SQL表来读取该主题:
CREATE TABLE Orders ( orderNumber STRING, orderDate STRING, productId STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'Orders', 'properties.bootstrap.servers' = 'xxxxxx', 'properties.group.id' = 'FlinkGroupId', 'format' = 'json', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'key.fields' = 'orderNumber;' );
问题现象
执行查询语句select * from Orders;时,orderNumber字段显示为NULL而非实际值1234:
- 尝试将
orderNumber字段类型改为STRING(原定义已是STRING),结果仍为NULL - 若发送格式错误的JSON作为Key,Flink会抛出解析错误,说明Flink确实在处理Key字段,但为何查询输出中该字段为
NULL?
问题原因与解决方案
问题出在两处配置错误:
key.fields字段名错误:配置项末尾多了一个分号,Flink会将目标字段识别为orderNumber;,但Kafka消息Key中的实际字段名是orderNumber,字段名不匹配导致无法映射值,最终显示为NULL。- 缺少
key.fields-include配置:默认情况下key.fields-include的值为ALL,意味着Flink会认为表中所有字段都来自Kafka消息的Key,但你的orderDate和productId实际存储在Value中,必须明确配置key.fields-include = 'EXCEPT_KEY',指定key.fields中定义的字段来自Key,其余字段来自Value。
修正后的表创建语句
CREATE TABLE Orders ( orderNumber STRING, orderDate STRING, productId STRING ) WITH ( 'connector' = 'kafka', 'topic' = 'Orders', 'properties.bootstrap.servers' = 'xxxxxx', 'properties.group.id' = 'FlinkGroupId', 'format' = 'json', 'scan.startup.mode' = 'earliest-offset', 'key.format' = 'json', 'key.fields' = 'orderNumber', 'key.fields-include' = 'EXCEPT_KEY' );
修正后重新执行查询,orderNumber字段就能正确读取到Key中的值1234了。
内容的提问来源于stack exchange,提问作者user3174411
相关产品推荐
相关产品推荐

