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

如何使用Flink SQL读取Kafka Topic的Key

要正确读取Kafka消息的Key,需在表定义中补充Key字段映射规则和必要的连接器配置,具体调整方案如下:

1. 核心配置修正

当前表定义缺少key.fields配置项——这是关联Kafka Key与表字段的关键参数,同时需要确保表结构与Key的字段对应:

修改后的完整表定义

CREATE OR REPLACE TABLE table /** mode('streaming')*/
(
        `pk` string,
        `id` string,
        -- 可选:声明主键,与Kafka Key映射关联
        PRIMARY KEY (`pk`) NOT ENFORCED
)WITH (
        'connector' = 'kafka',
        'key.format' = 'json',
        'key.fields' = 'pk', -- 指定表中`pk`字段对应Kafka Key的内容
        'key.fields-prefix' = '', -- 若Key字段与表字段无重名冲突,留空即可
        'properties.bootstrap.servers' = 'kafka:29092',
        'properties.group.id' = 'ef_entitlement_local_1',
        'topic-pattern' = 'entitlement_v1__.+?',
        'format' = 'json',
        'scan.startup.mode' = 'earliest-offset',
        'json.timestamp-format.standard' = 'ISO-8601',
        'json.fail-on-missing-field' = 'false',
        'json.ignore-parse-errors' = 'true'
);

2. 特殊场景适配

  • 如果Kafka Key是纯字符串/非JSON格式,直接替换key.format配置:
    'key.format' = 'raw',
    'key.format.derive-schema' = 'true'
    
  • 如果Key包含多字段结构,用逗号分隔字段名配置key.fields,比如'key.fields' = 'pk,id',同时确保表结构包含对应字段。

3. 验证配置有效性

执行以下查询确认Kafka Key是否正常读取:

SELECT pk, id FROM table LIMIT 10;

若pk字段返回Kafka Key的实际内容,说明配置生效。

内容的提问来源于stack exchange,提问作者hitesh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 00:12:08