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

Flink Upsert-Kafka输出Null值及Join报错问题求助

问题解决方案

一、移除upsert-kafka中的DELETE类型Null值

问题根源

左连接操作会产生回撤(DELETE)消息:当ROLES_DELETED_V1中的记录因状态清理、重复消费等原因被回撤时,左连接结果会生成DELETE类型的changelog,对应final_topic中的Null值。要保留最终状态,需将左连接结果转换为仅包含最新INSERT/UPDATE的消息。

解决步骤

  1. 补充缺失字段:原始ROLES_NORMALIZED视图缺少name、org、pod、event_timestamp等后续逻辑需要的字段,先修正视图:
CREATE VIEW ROLES_NORMALIZED AS
(
SELECT 
       JSON_VALUE(contentJson, '$.id') AS id,
       rr.type AS type,
       JSON_VALUE(contentJson, '$.name') AS name,
       JSON_VALUE(contentJson, '$.org') AS org,
       JSON_VALUE(contentJson, '$.pod') AS pod,
       TO_TIMESTAMP_LTZ(JSON_VALUE(contentJson, '$.modified')) AS modified,
       -- 从contentJson提取事件时间,或用当前时间,根据实际业务调整
       TO_TIMESTAMP_LTZ(JSON_VALUE(contentJson, '$.eventTimestamp')) AS event_timestamp
FROM raw_table rr
);
  1. 聚合保留最新状态:通过ROW_NUMBER()窗口函数过滤每个id的最新记录,将changelog转换为仅保留最终状态的消息:
INSERT INTO final_topic
SELECT
    event_timestamp,
    id,
    name,
    deleted,
    org,
    pod
FROM (
    SELECT
        GREATEST(r.event_timestamp, COALESCE(d.event_timestamp, r.event_timestamp)) AS event_timestamp,
        r.id,
        r.name,
        d.deleted,
        r.org,
        r.pod,
        -- 按事件时间倒序,取每个id的最新记录
        ROW_NUMBER() OVER (PARTITION BY r.id ORDER BY GREATEST(r.event_timestamp, COALESCE(d.event_timestamp, r.event_timestamp)) DESC) AS rn
    FROM ROLES_UPSERTS_V1 r
    LEFT JOIN ROLES_DELETED_V1 d ON r.id = d.id
) t
WHERE rn = 1;

此方式会确保每个id仅保留最新的状态,不会产生DELETE类型的回撤消息,也就不会出现Null值。


二、普通Kafka连接器支持左连接结果输出

问题根源

普通Kafka连接器是append-only模式,仅能接收INSERT类型的消息;而左连接会生成UPDATE/DELETE类型的changelog,因此直接写入会报错。需将changelog转换为append-only的流后再输出。

解决步骤

  1. 创建普通Kafka输出表:
CREATE TABLE final_topic_append
(
    op_type STRING, -- 标记操作类型:INSERT/UPDATE
    event_timestamp TIMESTAMP_LTZ,
    id VARCHAR,
    name VARCHAR,
    deleted TIMESTAMP_LTZ,
    org VARCHAR,
    pod VARCHAR
) WITH (
    'connector' = 'kafka',
    'topic' = 'final_topic_append',
    'properties.bootstrap.servers' = 'localhost:29092,localhost:39092',
    'format' = 'json',
    'scan.startup.mode' = 'earliest-offset',
    'properties.allow.auto.create.topics' = 'true'
);
  1. 转换changelog为append流:使用TO_CHANGELOG函数解析左连接后的changelog,过滤掉DELETE操作,仅保留INSERT/UPDATE消息:
INSERT INTO final_topic_append
SELECT
    op AS op_type,
    event_timestamp,
    id,
    name,
    deleted,
    org,
    pod
FROM TABLE(
    TO_CHANGELOG(
        SELECT
            GREATEST(r.event_timestamp, COALESCE(d.event_timestamp, r.event_timestamp)) AS event_timestamp,
            r.id,
            r.name,
            d.deleted,
            r.org,
            r.pod
        FROM ROLES_UPSERTS_V1 r
        LEFT JOIN ROLES_DELETED_V1 d ON r.id = d.id
    )
) t
WHERE op IN ('INSERT', 'UPDATE_AFTER');

也可以用之前的聚合方式直接生成最新状态,以UPDATE类型的append消息输出,效果一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 16:00:59