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

使用Trino处理Iceberg表:区分NULL类型并获取最新记录

Trino查询Iceberg表:按POS取最新列值并区分NULL类型

问题背景

通过Spark读取JSON文件创建的Iceberg表中存在两种NULL值:

  • 实际NULL:原始JSON中字段明确为NULL
  • 缺失生成的NULL:原始JSON中无该字段,Spark自动补全的NULL

需通过Trino编写查询,按pos字段(版本号,递增)排序获取各列最新有效值,同时借助changed_cols字段(数组类型,记录本次修改的字段)区分两种NULL——仅当字段在changed_cols中时,对应记录的NULL才是实际值,否则为缺失生成的NULL,需忽略并保留之前的有效值。

示例数据

假设Iceberg表user_data的历史数据如下:

idnameageaddressposchanged_cols
1Alice25NULL1["name","age","address"]
1NULLNULLBeijing2["address"]
1BobNULLNULL3["name"]

预期输出

idnameageaddress
1Bob25Beijing

解决方案SQL

核心思路

针对每个字段,仅考虑changed_cols包含该字段的记录(即有效修改记录),按pos降序排序后取第一条值,自动忽略缺失生成的NULL。

SELECT DISTINCT
    id,
    -- 获取name字段最新有效值
    FIRST_VALUE(name) OVER (
        PARTITION BY id 
        ORDER BY CASE WHEN array_contains(changed_cols, 'name') THEN pos ELSE -1 END DESC
        ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
    ) AS name,
    -- 获取age字段最新有效值
    FIRST_VALUE(age) OVER (
        PARTITION BY id 
        ORDER BY CASE WHEN array_contains(changed_cols, 'age') THEN pos ELSE -1 END DESC
        ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
    ) AS age,
    -- 获取address字段最新有效值
    FIRST_VALUE(address) OVER (
        PARTITION BY id 
        ORDER BY CASE WHEN array_contains(changed_cols, 'address') THEN pos ELSE -1 END DESC
        ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
    ) AS address,
    -- 可选:聚合去重所有被修改过的字段
    array_distinct(flatten(collect_list(changed_cols))) OVER (PARTITION BY id) AS all_changed_cols
FROM user_data;

代码解释

  1. 窗口排序逻辑:
    • 对每个字段,用CASE语句标记有效修改记录:若字段在changed_cols中,用实际pos值排序,否则标记为-1(排在最后)
    • 按标记后的"有效pos"降序排序,确保最新的有效修改记录排在最前
  2. FIRST_VALUE函数:取窗口内第一条记录的字段值,即为该字段的最新有效值
  3. 聚合去重changed_cols:用flatten+collect_list+array_distinct合并所有分组内的修改字段,得到该id所有被修改过的字段列表

测试用JSON数据(模拟Spark输入)

对应示例数据的原始JSON:

  1. {"id":1,"name":"Alice","age":25,"address":null} → 生成pos=1的记录,changed_cols包含所有字段
  2. {"id":1,"address":"Beijing"} → 生成pos=2的记录,changed_cols仅包含address
  3. {"id":1,"name":"Bob"} → 生成pos=3的记录,changed_cols仅包含name

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 06:09:59