You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink SQL Client读取Kafka消息Key时orderNumber字段为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表来读取该主题:

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?

问题原因与解决方案

问题出在两处配置错误:

  1. key.fields字段名错误:配置项末尾多了一个分号,Flink会将目标字段识别为orderNumber;,但Kafka消息Key中的实际字段名是orderNumber,字段名不匹配导致无法映射值,最终显示为NULL。
  2. 缺少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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 23:50:17