欢迎光临
我们一直在努力

Apache Iceberg 学习指南

Apache Iceberg 完整学习指南:从入门到进阶(2026 版)

适用版本:Apache Iceberg 1.11.0(2026-05-19 发布)| Format Version 3 最后更新:2026年8月 目标读者:数据工程师、平台工程师、数据架构师、对湖仓一体感兴趣的后端开发者


目录

入门篇

  • Apache Iceberg 全景概述
  • 核心概念与术语
  • 底层存储架构深度解析
  • 环境搭建与快速开始
  • DDL 操作详解
  • DML 操作详解
  • 查询与读取优化
  • PyIceberg Python SDK 入门
  • 进阶篇

  • Iceberg v3 新特性深度解读
  • Iceberg 与计算引擎集成
  • Catalog 选型与部署
  • 性能调优
  • 数据湖 CDC 与流批一体
  • 表维护与生命周期管理
  • 安全与权限
  • 生产环境部署架构
  • 实战项目
  • 2025-2026 技术趋势

  • 第一部分:入门篇


    1. Apache Iceberg 全景概述

    1.1 什么是 Apache Iceberg?

    Apache Iceberg 是一个开放的表格式(Open Table Format),为存储在数据湖中的大规模数据文件提供数据库级别的能力。你可以把它理解为:

    Iceberg = 数据湖的"数据库引擎"

    传统数据湖只是把文件(Parquet、ORC、CSV)堆在对象存储上,没有事务、没有 Schema 演进、没有高效的分区管理。Iceberg 在数据文件之上构建了一层元数据层,赋予数据湖以下能力:

    • ✅ ACID 事务:完全可序列化隔离,读写互不干扰
    • ✅ Schema Evolution:增删改列名,无需重写数据文件
    • ✅ 隐藏分区:用户无需感知分区列,查询自动裁剪
    • ✅ Time Travel:查询历史快照,回滚数据
    • ✅ Format Evolution:分区方案可以在线变更,无需重写数据
    • ✅ 多引擎兼容:Spark、Flink、Trino、Snowflake、Databricks 等均可读写同一张表

    1.2 诞生背景

    时间线:
    2017 ─── Netflix 内部研发,解决 Hive 表的性能与可靠性问题
    2018 ─── 贡献给 Apache Software Foundation,进入孵化期
    2020 ─── 成为 Apache 顶级项目(Top-Level Project)
    2022 ─── Iceberg v2 规范发布(行级删除、序列号)
    2023 ─── REST Catalog 规范确立,生态爆发
    2024 ─── Snowflake/Databricks/Google 全面支持
    2025 ─── Iceberg v3 规范批准(Deletion Vectors、Row Lineage、VARIANT)
    2026.05 ─── Apache Iceberg 1.11.0 发布,v3 特性全面落地

    Netflix 为什么需要 Iceberg?

    Netflix 在 2017 年面临的核心问题:

  • Hive 表不可靠:并发写入导致数据损坏,没有事务保证
  • 分区方案僵化:从日分区改为小时分区需要 PB 级数据重写
  • Schema 变更痛苦:新增一列需要重写所有数据文件
  • 小文件灾难:大量 Spark 作业产生百万级小文件,查询极慢
  • 缺乏时间旅行:数据错误无法回滚,只能手动修复
  • 这些痛点催生了 Iceberg——一个从底层重新设计的表格式。

    1.3 湖仓一体架构

    ┌──────────────────────────────────────────────────────────────┐
    │ 湖仓一体(Lakehouse)架构 │
    ├──────────────────────────────────────────────────────────────┤
    │ │
    │ ┌─────────┐ ┌──────────┐ ┌────────┐ ┌──────────────┐ │
    │ │ Spark │ │ Flink │ │ Trino │ │ Snowflake │ │
    │ │ (批处理) │ │ (流处理) │ │ (交互) │ │ (云数仓) │ │
    │ └────┬────┘ └────┬─────┘ └───┬────┘ └──────┬───────┘ │
    │ │ │ │ │ │
    │ └────────────┴─────┬──────┴───────────────┘ │
    │ │ │
    │ ┌───────────┴───────────┐ │
    │ │ Apache Iceberg │ ← 表格式层 │
    │ │ (ACID/Schema/分区) │ │
    │ └───────────┬───────────┘ │
    │ │ │
    │ ┌───────────┴───────────┐ │
    │ │ Parquet / ORC 文件 │ ← 数据文件层 │
    │ └───────────┬───────────┘ │
    │ │ │
    │ ┌───────────────────────┴───────────────────────────┐ │
    │ │ 对象存储 (S3 / HDFS / MinIO / ADLS) │ │
    │ └───────────────────────────────────────────────────┘ │
    │ │
    └──────────────────────────────────────────────────────────────┘

    1.4 传统数仓 vs 数据湖 vs 湖仓一体

    维度传统数据仓库传统数据湖湖仓一体(Iceberg)
    存储 专有格式,绑定硬件 开放文件(Parquet) 开放文件 + 表格式
    事务 完整 ACID 无事务 完整 ACID
    Schema 演进 有限支持 需要重写 在线演进
    分区 静态分区 静态分区 隐藏分区 + 分区演进
    时间旅行 不支持 不支持 原生支持
    计算引擎 单一引擎 多引擎(但无一致性) 多引擎(ACID 保证)
    存储成本 高(专有存储) 低(对象存储) 低(对象存储)
    数据质量 低(可能脏数据) 高(ACID 保证)

    1.5 四大开源表格式对比

    特性Apache IcebergDelta LakeApache HudiApache Paimon
    起源 Netflix (2017) Databricks (2019) Uber (2019) Alibaba Flink (2022)
    治理 Apache 基金会 Apache 基金会 Apache 基金会 Apache 基金会
    ACID 事务 ✅ 完整序列化隔离 ✅ 完整 ✅ 完整 ✅ 完整
    Schema Evolution ✅ 增删改名重排 ✅ 增删改名 ✅ 增删改 ✅ 增删改
    分区演进 ✅ 无需重写 ⚠️ 有限 ✅ 支持 ✅ 支持
    Time Travel ✅ 快照查询 ✅ 快照查询 ✅ 支持 ✅ 支持
    行级删除 ✅ v2/v3 Deletion Vectors ✅ Deletion Vectors ✅ 原生索引 ✅ LSM-tree
    CDC/流式写入 ✅ Flink/MOR ⚠️ Structured Streaming ✅ 强项 ✅ 强项
    多引擎支持 ✅ 最广 ⚠️ 偏 Spark/Databricks ⚠️ Spark 为主 ⚠️ Flink 为主
    Catalog 生态 REST/Hive/JDBC/Glue/Nessie Unity Catalog Hive/Hadoop Hive/Filesystem
    最佳场景 通用湖仓、多引擎 Databricks 生态 高频 Upsert Flink 流式优先
    2026 生态位 事实标准 Databricks 深度绑定 CDC 特化 Flink 特化

    🏆 面试加分提示:2026 年 Iceberg 已经成为事实上的标准表格式。Iceberg Summit 2026 旧金山站吸引超过 600 名参会者、70+ 场次技术分享,来自 Google、Apple、Snowflake、Databricks、Microsoft、Netflix、LinkedIn 的工程师共同参与设计讨论。这种跨厂商的参与度才是 Iceberg 真正的护城河。

    ⚠️ 生产踩坑提醒:选择表格式不要只看功能列表,要确认你使用的计算引擎有成熟的原生连接器。技术 superior 的格式如果没有你主力引擎的成熟 connector,风险比收益大。


    2. 核心概念与术语

    2.1 五层元数据结构

    理解 Iceberg 的核心是理解它的五层元数据结构。这是 Iceberg 设计最精妙的部分:

    ┌─────────────────────────────────────────────────────────────────┐
    │ Iceberg 五层元数据结构 │
    ├─────────────────────────────────────────────────────────────────┤
    │ │
    │ Layer 1: Metadata File (metadata.json) │
    │ ├── 表的完整状态:schema、分区规范、排序顺序 │
    │ ├── 当前快照 ID、快照列表 │
    │ └── 表属性配置 │
    │ │
    │ Layer 2: Snapshot (快照) │
    │ ├── 表在某一时刻的完整视图 │
    │ ├── 指向一个 Manifest List │
    │ └── 每个操作(INSERT/DELETE)产生新快照 │
    │ │
    │ Layer 3: Manifest List (snap-*.avro) │
    │ ├── 列出该快照包含的所有 Manifest File │
    │ ├── 记录每个 Manifest 的分区统计信息 │
    │ └── 用于快速跳过不相关的 Manifest │
    │ │
    │ Layer 4: Manifest File (*.avro) │
    │ ├── 列出该 Manifest 包含的所有 Data File │
    │ ├── 每个 Data File 的分区值、记录数、列统计(min/max/null_count) │
    │ └── 记录文件状态:ADDED / EXISTING / DELETED │
    │ │
    │ Layer 5: Data File (Parquet/ORC) │
    │ ├── 实际存储数据的文件 │
    │ ├── 列式存储,自带压缩 │
    │ └── 默认格式:Parquet(v3 推荐使用 ZSTD 压缩) │
    │ │
    └─────────────────────────────────────────────────────────────────┘

    2.2 关键术语详解

    Table(表)

    Iceberg 中的一张表 = 元数据文件 + 数据文件。表的定义包括:

    • Schema:列名、类型、唯一 ID
    • Partition Spec:分区转换规则
    • Sort Order:数据排序方式
    • Properties:表属性配置
    • Snapshots:历史快照集合
    • Current Snapshot:当前有效数据
    Snapshot(快照)

    每次写入操作产生一个新快照。快照是 Iceberg 实现 ACID 和时间旅行的基础:

    时间线示例:

    t0 ─── 初始快照 (snapshot-0)
    │ └── 包含 10 个数据文件,100 万行

    t1 ─── INSERT 操作 (snapshot-1)
    │ └── 新增 2 个数据文件,20 万行
    │ └── 总数据量:12 个文件,120 万行

    t2 ─── DELETE 操作 (snapshot-2)
    │ └── 删除 1 个数据文件中的部分行
    │ └── 总数据量:11 个文件 + 1 个删除文件

    t3 ─── 查询 snapshot-0 → 看到 t0 时刻的数据(100 万行)
    t4 ─── 查询 snapshot-2 → 看到最新数据

    Catalog(目录服务)

    Catalog 是 Iceberg 表的注册中心,负责:

    • 存储表名到元数据文件位置的映射
    • 协调并发写入(原子提交)
    • 管理命名空间(namespace)

    常见 Catalog 类型:Hive、REST、JDBC、Nessie、Glue

    Format Version(格式版本)
    版本发布关键特性
    v1 2018 基础 ACID、Schema Evolution、隐藏分区
    v2 2022 行级删除(Position Delete / Equality Delete)、序列号
    v3 2025-2026 Deletion Vectors、Row Lineage、VARIANT 类型、纳秒时间戳、Geometry 类型、默认列值

    🏆 面试加分提示:Iceberg 的版本号容易混淆。库版本(如 1.11.0)是 Java 实现的版本号;格式版本(如 v3)是磁盘上的契约规范,任何语言的任何实现都必须遵守。1.11.0 是当前稳定的库版本,全面支持 v3 格式。

    ⚠️ 生产踩坑提醒:升级到 v3 前,确保所有读写引擎都支持 v3。如果有任何引擎不兼容,保持 v2 格式,否则会导致数据一致性问题。


    3. 底层存储架构深度解析

    3.1 物理存储布局

    warehouse/
    └── db/
    └── table_name/
    ├── metadata/
    │ ├── 00000-00000000-0000-0000-0000-000000000000.metadata.json
    │ ├── 00001-00000000-0000-0000-0000-000000000001.metadata.json
    │ ├── snap-1089405398375037-0-xxxx.avro ← Manifest List
    │ ├── snap-1089405398375037-1-xxxx.avro ← Manifest List
    │ ├── xxxxxxxx-0000-0000-0000-xxxxxxxxxxxx-m0.avro ← Manifest File
    │ └── xxxxxxxx-0000-0000-0000-xxxxxxxxxxxx-m1.avro ← Manifest File
    ├── data/
    │ ├── partition_col=2024-01-01/
    │ │ ├── 00000-0-xxxx-00001.parquet ← Data File
    │ │ └── 00001-0-xxxx-00002.parquet
    │ └── partition_col=2024-01-02/
    │ └── 00000-0-xxxx-00003.parquet
    └── (v2/v3) delete files or deletion vectors

    3.2 Metadata JSON 文件结构

    {
    "format-version": 3,
    "table-uuid": "a1b2c3d4-…",
    "location": "s3://warehouse/db/table",
    "last-sequence-number": 42,
    "last-updated-ms": 1716000000000,
    "last-column-id": 10,
    "schemas": [
    {
    "type": "struct",
    "schema-id": 0,
    "fields": [
    {"id": 1, "name": "id", "required": true, "type": "long"},
    {"id": 2, "name": "name", "required": false, "type": "string"},
    {"id": 3, "name": "created_at", "required": false, "type": "timestamptz"}
    ]
    }
    ],
    "current-schema-id": 0,
    "partition-specs": [
    {
    "spec-id": 0,
    "fields": [
    {"source-id": 3, "field-id": 1000, "name": "created_at_day", "transform": "day"}
    ]
    }
    ],
    "default-spec-id": 0,
    "sort-orders": [
    {"order-id": 1, "fields": [{"transform": "identity", "source-id": 1, "direction": "asc", "null-order": "nulls-first"}]}
    ],
    "default-sort-order-id": 1,
    "properties": {
    "format-version": "3",
    "write.delete.mode": "merge-on-read"
    },
    "current-snapshot-id": 1089405398375037,
    "snapshots": [
    {
    "snapshot-id": 1089405398375037,
    "timestamp-ms": 1716000000000,
    "summary": {"operation": "append", "added-data-files": "2"},
    "manifest-list": "s3://warehouse/db/table/metadata/snap-1089405398375037-0.avro",
    "schema-id": 0
    }
    ],
    "refs": {
    "main": {"snapshot-id": 1089405398375037, "type": "branch"}
    },
    "next-row-id": 500000
    }

    3.3 隐藏分区(Hidden Partitioning)

    传统分区需要用户在 WHERE 子句中显式指定分区列:

    — 传统 Hive 分区:用户需要知道分区结构
    SELECT * FROM events WHERE date_hour >= '2024-01-01-00' AND date_hour < '2024-01-02-00';

    — Iceberg 隐藏分区:用户按原始列过滤,引擎自动推导分区裁剪
    SELECT * FROM events WHERE event_time >= '2024-01-01 00:00:00' AND event_time < '2024-01-02 00:00:00';

    3.4 分区转换函数

    转换函数说明示例
    identity 原始值作为分区值 identity(col)
    year 提取年份 year(timestamp_col) → 2024
    month 提取年月 month(timestamp_col) → 2024-01
    day 提取日期 day(timestamp_col) → 2024-01-15
    hour 提取小时 hour(timestamp_col) → 2024-01-15-08
    bucket(N) 哈希分桶 bucket(64, device_id) → 0~63
    truncate(W) 前缀截断 truncate(10, name) → 前 10 个字符
    void 虚拟分区(不分区) void(col)

    3.5 Manifest File 详解

    Manifest File 是 Avro 格式文件,每个条目包含:

    Data File Entry:
    ├── file_path: 数据文件的完整路径
    ├── file_format: PARQUET / ORC / AVRO
    ├── partition: 分区值(结构体)
    ├── record_count: 记录数
    ├── file_size_in_bytes: 文件大小
    ├── column_sizes: 各列大小(Map<column_id, size>)
    ├── value_counts: 各列非空值数量
    ├── null_value_counts: 各列空值数量
    ├── nan_value_counts: 各列 NaN 值数量
    ├── lower_bounds: 各列最小值(Map<column_id, value>)
    ├── upper_bounds: 各列最大值(Map<column_id, value>)
    ├── split_offsets: 文件的分割偏移量列表
    ├── sort_order_id: 排序顺序 ID
    ├── status: ADDED / EXISTING / DELETED
    ├── spec_id: 分区规范 ID
    └── (v3) first_row_id: 文件首个行 ID(Row Lineage 支持)

    🏆 面试加分提示:Manifest 文件中的 lower_bounds 和 upper_bounds 是实现 Filter Pushdown(过滤条件下推)的关键。引擎无需打开数据文件,仅通过 Manifest 的列统计信息就能跳过大量不相关的数据文件,这就是 Iceberg 查询性能远超 Hive 的核心原因之一。

    ⚠️ 生产踩坑提醒:Metadata 目录会随着快照累积产生大量文件。务必配置 write.metadata.delete-after-commit.enabled = true 和 history.expire.max-snapshot-age-ms 来自动清理旧元数据,否则 Metadata 目录可能累积数十万个文件,严重影响性能。


    4. 环境搭建与快速开始

    4.1 Docker Compose 搭建本地环境

    这是最快的入门方式,一键启动完整的 Iceberg 开发环境:

    # docker-compose.yml
    services:
    spark-iceberg:
    image: tabulario/sparkiceberg
    container_name: sparkiceberg
    networks:
    iceberg_net:
    depends_on:
    rest
    minio
    volumes:
    ./warehouse:/home/iceberg/warehouse
    ./notebooks:/home/iceberg/notebooks/notebooks
    environment:
    AWS_ACCESS_KEY_ID=admin
    AWS_SECRET_ACCESS_KEY=password
    AWS_REGION=useast1
    ports:
    8888:8888 # Jupyter Notebook
    8080:8080 # Spark Master UI
    10000:10000 # Spark Thrift Server
    10001:10001

    rest:
    image: apache/icebergrestfixture
    container_name: icebergrest
    networks:
    iceberg_net:
    ports:
    8181:8181
    environment:
    AWS_ACCESS_KEY_ID=admin
    AWS_SECRET_ACCESS_KEY=password
    AWS_REGION=useast1
    CATALOG_WAREHOUSE=s3://warehouse/
    CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO
    CATALOG_S3_ENDPOINT=http://minio:9000

    minio:
    image: minio/minio
    container_name: minio
    environment:
    MINIO_ROOT_USER=admin
    MINIO_ROOT_PASSWORD=password
    MINIO_DOMAIN=minio
    networks:
    iceberg_net:
    aliases:
    warehouse.minio
    ports:
    9001:9001 # MinIO Console
    9000:9000 # S3 API
    command: ["server", "/data", "–console-address", ":9001"]

    mc:
    depends_on:
    minio
    image: minio/mc
    container_name: mc
    networks:
    iceberg_net:
    environment:
    AWS_ACCESS_KEY_ID=admin
    AWS_SECRET_ACCESS_KEY=password
    AWS_REGION=useast1
    entrypoint: |
    /bin/sh -c "
    until (/usr/bin/mc alias set minio http://minio:9000 admin password) do
    echo '…waiting…' && sleep 1;
    done;
    /usr/bin/mc rm -r –force minio/warehouse;
    /usr/bin/mc mb minio/warehouse;
    /usr/bin/mc policy set public minio/warehouse;
    tail -f /dev/null"

    networks:
    iceberg_net:

    启动命令:

    docker-compose up -d

    4.2 PySpark 连接 Iceberg

    from pyspark.sql import SparkSession

    # 创建 Spark Session,配置 Iceberg 连接
    spark = SparkSession.builder \\
    .appName("IcebergQuickStart") \\
    .config("spark.sql.extensions",
    "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \\
    .config("spark.sql.catalog.demo", "org.apache.iceberg.spark.SparkCatalog") \\
    .config("spark.sql.catalog.demo.type", "rest") \\
    .config("spark.sql.catalog.demo.uri", "http://localhost:8181") \\
    .config("spark.sql.catalog.demo.warehouse", "s3://warehouse") \\
    .config("spark.sql.catalog.demo.io-impl",
    "org.apache.iceberg.aws.s3.S3FileIO") \\
    .config("spark.sql.catalog.demo.s3.endpoint", "http://localhost:9000") \\
    .config("spark.sql.catalog.demo.s3.path-style-access", "true") \\
    .config("spark.hadoop.fs.s3a.access.key", "admin") \\
    .config("spark.hadoop.fs.s3a.secret.key", "password") \\
    .getOrCreate()

    # 创建命名空间
    spark.sql("CREATE NAMESPACE IF NOT EXISTS demo.analytics")

    # 创建表
    spark.sql("""
    CREATE TABLE IF NOT EXISTS demo.analytics.events (
    event_id BIGINT,
    user_id BIGINT,
    event_type STRING,
    payload STRING,
    event_time TIMESTAMP
    ) USING iceberg
    PARTITIONED BY (days(event_time))
    TBLPROPERTIES (
    'format-version' = '3',
    'write.parquet.compression-codec' = 'zstd'
    )
    """
    )

    # 插入数据
    spark.sql("""
    INSERT INTO demo.analytics.events VALUES
    (1, 1001, 'click', '{"page": "/home"}', timestamp('2024-01-15 10:00:00')),
    (2, 1002, 'purchase', '{"item": "book"}', timestamp('2024-01-15 11:00:00')),
    (3, 1001, 'click', '{"page": "/about"}', timestamp('2024-01-15 12:00:00'))
    """
    )

    # 查询数据
    spark.sql("SELECT * FROM demo.analytics.events").show()

    # 查看表结构
    spark.sql("DESCRIBE TABLE demo.analytics.events").show()

    4.3 Spark SQL CLI 快速连接

    spark-sql \\
    –packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0 \\
    –conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \\
    –conf spark.sql.catalog.demo=org.apache.iceberg.spark.SparkCatalog \\
    –conf spark.sql.catalog.demo.type=rest \\
    –conf spark.sql.catalog.demo.uri=http://localhost:8181 \\
    –conf spark.sql.catalog.demo.warehouse=s3://warehouse \\
    –conf spark.sql.catalog.demo.io-impl=org.apache.iceberg.aws.s3.S3FileIO \\
    –conf spark.sql.catalog.demo.s3.endpoint=http://localhost:9000 \\
    –conf spark.sql.catalog.demo.s3.path-style-access=true

    4.4 REST Catalog 独立部署

    生产环境推荐使用 PostgreSQL 作为 REST Catalog 的后端存储:

    # 使用 Apache Polaris(2026年2月毕业为 Apache 顶级项目)
    docker run -d \\
    –name polaris \\
    -p 8181:8181 \\
    -e POLARIS_DB_URL=jdbc:postgresql://postgres:5432/polaris \\
    -e POLARIS_DB_USER=polaris \\
    -e POLARIS_DB_PASSWORD=polaris \\
    apache/polaris:latest

    # 或使用轻量级 Lakekeeper(Rust 实现)
    docker run -d \\
    –name lakekeeper \\
    -p 8181:8181 \\
    -e LAKEKEEPER__PG_DATABASE_URL_READ=postgres://catalog:catalog@postgres:5432/catalog \\
    -e LAKEKEEPER__PG_DATABASE_URL_WRITE=postgres://catalog:catalog@postgres:5432/catalog \\
    -e LAKEKEEPER__BASE_URI=http://localhost:8181 \\
    quay.io/lakekeeper/catalog:latest serve

    ⚠️ 生产踩坑提醒:本地开发用 tabulario/spark-iceberg 镜像很方便,但生产环境不要用。生产环境应该使用独立的 REST Catalog(如 Polaris、Lakekeeper)配合对象存储。


    5. DDL 操作详解

    5.1 CREATE TABLE

    — 基础建表
    CREATE TABLE demo.analytics.orders (
    order_id BIGINT,
    customer_id BIGINT,
    order_date DATE,
    total_amount DECIMAL(10, 2),
    status STRING,
    created_at TIMESTAMP,
    updated_at TIMESTAMP
    ) USING iceberg
    TBLPROPERTIES (
    'format-version' = '3',
    'write.format.default' = 'parquet',
    'write.parquet.compression-codec' = 'zstd',
    'commit.retry.num-workers' = '4'
    );

    — 带分区和排序的建表
    CREATE TABLE demo.analytics.page_views (
    view_id BIGINT,
    user_id BIGINT,
    page_url STRING,
    referrer_url STRING,
    user_agent STRING,
    view_time TIMESTAMP,
    ip_address STRING
    ) USING iceberg
    PARTITIONED BY (
    days(view_time), — 按天分区
    bucket(64, user_id) — 按 user_id 哈希 64 桶
    )
    TBLPROPERTIES (
    'format-version' = '3',
    'write.parquet.compression-codec' = 'zstd',
    'write.target-file-size-bytes' = '536870912',
    'write.metadata.delete-after-commit.enabled' = 'true',
    'write.metadata.previous-versions-max' = '50'
    );

    — 使用 CTAS(Create Table As Select)
    CREATE TABLE demo.analytics.user_summary USING iceberg AS
    SELECT
    user_id,
    COUNT(*) AS total_views,
    MAX(view_time) AS last_view_time
    FROM demo.analytics.page_views
    GROUP BY user_id;

    5.2 表属性配置参考

    属性默认值说明
    format-version 2 格式版本(1/2/3)
    write.format.default parquet 数据文件格式
    write.parquet.compression-codec zstd Parquet 压缩算法
    write.target-file-size-bytes 536870912 (512MB) 目标数据文件大小
    commit.retry.num-workers 1 提交重试并发数
    write.delete.mode copy-on-write (v1) / merge-on-read (v2+) 删除模式
    write.update.mode copy-on-write (v1) / merge-on-read (v2+) 更新模式
    write.merge.mode copy-on-write (v1) / merge-on-read (v2+) MERGE 模式
    history.expire.max-snapshot-age-ms 432000000 (5天) 快照最大保留时间
    history.expire.min-snapshots-to-keep 1 最少保留快照数
    write.metadata.delete-after-commit.enabled false 提交后删除旧元数据
    write.metadata.previous-versions-max 100 保留的最大旧版本数
    read.split.target-size 134217728 (128MB) 读取时目标分片大小
    read.split.planning-lookahead 10 文件规划前瞻数

    5.3 Schema Evolution(Schema 演进)

    Iceberg 的 Schema Evolution 是元数据级别操作,不需要重写数据文件:

    — 添加列(追加到末尾)
    ALTER TABLE demo.analytics.orders ADD COLUMNS (
    shipping_address STRING,
    discount_rate DECIMAL(5, 2)
    );

    — 添加列到指定位置
    ALTER TABLE demo.analytics.orders ADD COLUMNS (
    notes STRING AFTER status
    );

    — 重命名列
    ALTER TABLE demo.analytics.orders RENAME COLUMN notes TO order_notes;

    — 删除列
    ALTER TABLE demo.analytics.orders DROP COLUMN order_notes;

    — 修改列类型(仅支持兼容的类型拓宽)
    — 例如:INT → BIGINT, FLOAT → DOUBLE, DECIMAL(10,2) → DECIMAL(12,4)
    ALTER TABLE demo.analytics.orders ALTER COLUMN total_amount TYPE DECIMAL(12, 4);

    — 设置列默认值(v3 特性!对历史行透明应用)
    ALTER TABLE demo.analytics.orders ALTER COLUMN status SET DEFAULT 'pending';

    — 设置列为 NOT NULL(v3 特性:有默认值时可以添加 NOT NULL 列)
    ALTER TABLE demo.analytics.orders ALTER COLUMN status SET NOT NULL;

    Schema Evolution 支持的类型拓宽规则:

    原类型可拓宽为
    INT BIGINT
    FLOAT DOUBLE
    DECIMAL(P,S) DECIMAL(P’,S’) 其中 P’>=P, S’>=S
    —— ❌ 不支持 STRING ↔ INT 等不兼容转换

    5.4 Partition Evolution(分区演进)

    这是 Iceberg 的杀手级特性之一——在线更改分区方案,无需重写数据:

    — 原始分区:按天分区
    CREATE TABLE demo.analytics.events (...)
    USING iceberg
    PARTITIONED BY (days(event_time));

    — 改为按小时分区(只影响新写入的数据!)
    ALTER TABLE demo.analytics.events REPLACE PARTITION FIELD days(event_time) WITH hours(event_time);

    — 添加额外的分区维度
    ALTER TABLE demo.analytics.events ADD PARTITION FIELD bucket(16, user_id);

    — 移除分区字段
    ALTER TABLE demo.analytics.events DROP PARTITION FIELD days(event_time);

    分区演进示意:

    旧数据(按天分区) 新数据(按小时分区)
    ┌─────────────────┐ ┌─────────────────┐
    │ event_time_day │ │ event_time_hour │
    │ = 2024-01-15 │ │ = 2024-01-15-08 │
    │ ├── file1 │ │ ├── file5 │
    │ └── file2 │ │ └── file6 │
    └─────────────────┘ └─────────────────┘

    查询引擎自动判断使用哪个分区规范进行裁剪

    5.5 Sort Order 定义

    排序顺序影响查询性能,特别是与合并操作配合时:

    — 创建时定义排序顺序
    CREATE TABLE demo.analytics.orders (...)
    USING iceberg
    ORDERED BY (customer_id ASC, order_date DESC);

    — 修改排序顺序
    ALTER TABLE demo.analytics.orders SET ORDER BY (order_date DESC NULLS LAST);

    — 移除排序顺序
    ALTER TABLE demo.analytics.orders SET ORDER BY ();

    5.6 DROP TABLE 与 命名空间管理

    — 删除表(不可恢复!数据文件将被标记删除)
    DROP TABLE IF EXISTS demo.analytics.orders;

    — 创建命名空间
    CREATE NAMESPACE IF NOT EXISTS demo.analytics;

    — 列出命名空间
    SHOW NAMESPACES IN demo;

    — 删除命名空间(必须先删除其中的所有表)
    DROP NAMESPACE IF EXISTS demo.analytics;

    🏆 面试加分提示:Iceberg 的 Schema Evolution 基于列 ID而非列位置。每列在 Schema 中有唯一 ID,重命名列不影响 ID,因此历史查询仍然正确。这与 Hive 基于列位置的 Schema 管理形成鲜明对比。

    ⚠️ 生产踩坑提醒:分区演进后,查询计划器需要同时处理新旧两套分区规范。在过渡期内,某些查询可能无法完全利用分区裁剪。建议分区演进后逐步重写旧数据到新分区方案。


    6. DML 操作详解

    6.1 INSERT

    — 单行插入
    INSERT INTO demo.analytics.orders VALUES
    (1001, 5001, '2024-01-15', 299.99, 'completed', '2024-01-15 10:00:00', '2024-01-15 10:00:00');

    — 批量插入
    INSERT INTO demo.analytics.orders VALUES
    (1002, 5002, '2024-01-15', 159.00, 'pending', '2024-01-15 11:00:00', '2024-01-15 11:00:00'),
    (1003, 5001, '2024-01-16', 599.99, 'shipped', '2024-01-16 09:00:00', '2024-01-16 09:30:00');

    — 从其他表插入(CTAS / INSERT … SELECT)
    INSERT INTO demo.analytics.orders
    SELECT * FROM demo.staging.orders_import WHERE import_date = '2024-01-15';

    6.2 INSERT OVERWRITE

    — 覆盖整个表
    INSERT OVERWRITE demo.analytics.orders
    SELECT * FROM demo.staging.orders_full_refresh;

    — 覆盖特定分区
    INSERT OVERWRITE demo.analytics.orders
    PARTITION (order_date = '2024-01-15')
    SELECT order_id, customer_id, total_amount, status, created_at, updated_at
    FROM demo.staging.orders_daily
    WHERE order_date = '2024-01-15';

    — 动态分区覆盖
    INSERT OVERWRITE demo.analytics.orders
    SELECT order_id, customer_id, order_date, total_amount, status, created_at, updated_at
    FROM demo.staging.orders_full_refresh;

    6.3 Copy-on-Write vs Merge-on-Read

    这是 Iceberg 最核心的设计决策之一:

    ┌─────────────────────────────────────────────────────────────────┐
    │ Copy-on-Write (COW) vs Merge-on-Read (MOR) │
    ├─────────────────────────────────────────────────────────────────┤
    │ │
    │ Copy-on-Write (COW): │
    │ ┌──────────┐ DELETE/UPDATE ┌──────────┐ │
    │ │ File A │ ──────────────────→ │ File A' │ │
    │ │ (10 rows)│ 重写整个文件 │ (8 rows) │ │
    │ └──────────┘ └──────────┘ │
    │ 优点:读取快(无需合并) 缺点:写入慢(文件重写) │
    │ │
    │ Merge-on-Read (MOR): │
    │ ┌──────────┐ ┌──────────┐ │
    │ │ File A │ ── 保留原文件 ──→ │ File A │ │
    │ │ (10 rows)│ │ (10 rows)│ │
    │ └──────────┘ └────┬─────┘ │
    │ │ 叠加 │
    │ ┌────┴─────┐ │
    │ │ Delete │ │
    │ │ Vector │ ← v3 位图 │
    │ │ (2 rows) │ │
    │ └──────────┘ │
    │ 优点:写入快(不重写文件) 缺点:读取需要合并 │
    │ │
    └─────────────────────────────────────────────────────────────────┘

    维度Copy-on-Write (COW)Merge-on-Read (MOR)
    写入性能 慢(重写数据文件) 快(只写删除标记)
    读取性能 快(无需合并) 较慢(需合并删除标记)
    适用场景 读多写少、BI 报表 写多读少、CDC 入湖
    v3 中的删除标记 N/A Deletion Vector(位图)
    存储开销 高(文件副本) 低(位图很小)

    — 配置写入模式
    ALTER TABLE demo.analytics.orders SET TBLPROPERTIES (
    'write.delete.mode' = 'merge-on-read',
    'write.update.mode' = 'merge-on-read',
    'write.merge.mode' = 'merge-on-read'
    );

    6.4 DELETE

    — 按条件删除行
    DELETE FROM demo.analytics.orders WHERE order_id = 1001;

    — 批量删除
    DELETE FROM demo.analytics.orders WHERE status = 'cancelled' AND created_at < '2024-01-01';

    — 子查询删除
    DELETE FROM demo.analytics.orders
    WHERE customer_id IN (
    SELECT customer_id FROM demo.analytics.blocked_users
    );

    6.5 UPDATE

    — 更新单列
    UPDATE demo.analytics.orders SET status = 'refunded' WHERE order_id = 1002;

    — 更新多列
    UPDATE demo.analytics.orders
    SET status = 'shipped',
    updated_at = current_timestamp()
    WHERE status = 'processing' AND order_date >= '2024-01-15';

    6.6 MERGE INTO(Upsert 操作)

    — 经典的 Upsert 操作
    MERGE INTO demo.analytics.customers AS target
    USING demo.staging.customer_updates AS source
    ON target.customer_id = source.customer_id
    WHEN MATCHED AND source.is_deleted = true THEN
    DELETE
    WHEN MATCHED THEN
    UPDATE SET
    target.name = source.name,
    target.email = source.email,
    target.updated_at = current_timestamp()
    WHEN NOT MATCHED THEN
    INSERT (customer_id, name, email, created_at, updated_at)
    VALUES (source.customer_id, source.name, source.email,
    current_timestamp(), current_timestamp());

    — 增量合并(仅处理最近的数据)
    MERGE INTO demo.analytics.orders AS target
    USING (
    SELECT *, ROW_NUMBER() OVER (
    PARTITION BY order_id ORDER BY updated_at DESC
    ) AS rn
    FROM demo.staging.orders_cdc
    WHERE _batch_date = '2024-01-15'
    ) AS source
    ON target.order_id = source.order_id AND source.rn = 1
    WHEN MATCHED THEN
    UPDATE SET *
    WHEN NOT MATCHED THEN
    INSERT *;

    6.7 CDC 数据入湖模式

    ┌─────────────────────────────────────────────────────────────┐
    │ CDC 数据入湖架构 │
    ├─────────────────────────────────────────────────────────────┤
    │ │
    │ ┌─────────┐ ┌──────────┐ ┌────────┐ ┌─────────┐ │
    │ │ MySQL/ │───→│ Debezium │───→│ Kafka │───→│ Flink │ │
    │ │ Postgres│ │ (CDC) │ │ Topic │ │ Job │ │
    │ └─────────┘ └──────────┘ └────────┘ └────┬────┘ │
    │ │ │
    │ MERGE INTO │
    │ │ │
    │ ┌─────┴────┐ │
    │ │ Iceberg │ │
    │ │ Table │ │
    │ └──────────┘ │
    │ │
    └─────────────────────────────────────────────────────────────┘

    🏆 面试加分提示:v3 的 Deletion Vectors 将 DML 性能提升了约 10 倍。v2 中每次 MOR 删除需要创建单独的位置删除文件,读取时合并成本线性增长;v3 用一个紧凑的 Roaring Bitmap 位图替代,读取时只需做一次位图检查。这对 CDC 流式入湖场景是变革性的改进。

    ⚠️ 生产踩坑提醒:MERGE INTO 是 Iceberg 中最昂贵的操作。如果可能,用 INSERT OVERWRITE 覆盖整个分区代替逐行 MERGE。只有在需要保留同一分区内的历史数据时才使用 MERGE。


    7. 查询与读取优化

    7.1 Time Travel(时间旅行)

    — 查看表的快照历史
    SELECT * FROM demo.analytics.orders.snapshots;

    — 按快照 ID 查询
    SELECT * FROM demo.analytics.orders
    VERSION AS OF 1089405398375037;

    — 按时间戳查询(查询该时间点最近的快照)
    SELECT * FROM demo.analytics.orders
    TIMESTAMP AS OF '2024-01-15 12:00:00';

    — 按分支/标签查询
    SELECT * FROM demo.analytics.orders
    VERSION AS OF 'audit_branch';

    — 查看所有快照的详细信息
    SELECT
    snapshot_id,
    parent_id,
    operation,
    summary,
    committed_at
    FROM demo.analytics.orders.snapshots
    ORDER BY committed_at DESC;

    7.2 分支与标签(Branch/Tag)

    — 创建分支(基于当前快照)
    ALTER TABLE demo.analytics.orders CREATE BRANCH experiment_v2;

    — 创建标签(固定到特定快照)
    ALTER TABLE demo.analytics.orders CREATE TAG release_v1_0;

    — 基于分支写入数据(写入不影响 main 分支)
    — Spark:
    — df.writeTo("demo.analytics.orders").option("branch", "experiment_v2").append()

    — 切换分支进行查询
    SELECT * FROM demo.analytics.orders VERSION AS OF 'experiment_v2';

    — 将分支合并回 main
    ALTER TABLE demo.analytics.orders REPLACE BRANCH main FROM experiment_v2;

    — 删除分支/标签
    ALTER TABLE demo.analytics.orders DROP BRANCH experiment_v2;
    ALTER TABLE demo.analytics.orders DROP TAG release_v1_0;

    — 查看所有引用(分支和标签)
    SELECT * FROM demo.analytics.orders.refs;

    分支工作流示意:

    main ──── snap1 ──── snap2 ──── snap3 ──── snap4 (合并后)
    │ ↑
    └─ experiment ── s3a─s3b ─┘
    (独立开发,不影响 main)

    7.3 增量读取

    — 读取两个快照之间的增量数据
    — 使用 PyIceberg(推荐方式)
    — 或者使用 Spark 的 incremental read API

    — Spark SQL 查看增量数据(通过比较快照)
    SELECT * FROM demo.analytics.orders
    WHERE _change_type IS NOT NULL; — v3 行变更类型

    # PyIceberg 增量读取
    from pyiceberg.catalog import load_catalog

    catalog = load_catalog("demo")
    table = catalog.load_table("analytics.orders")

    # 读取快照 1 到快照 2 之间的增量
    incremental_df = table.scan(
    snapshot_id=1089405398375037,
    options={"start-snapshot-id": "1089405398375030"}
    ).to_arrow()

    7.4 过滤条件下推(Filter Pushdown)

    Iceberg 的三级数据跳过机制:

    ┌─────────────────────────────────────────────────────────────┐
    │ Iceberg 三级数据跳过机制 │
    ├─────────────────────────────────────────────────────────────┤
    │ │
    │ Level 1: Partition Pruning (分区裁剪) │
    │ ├── 基于隐藏分区的转换函数推导 │
    │ ├── 查询 WHERE event_time > '2024-01-15' │
    │ └── 自动跳过 event_time < 2024-01-15 的整个分区 │
    │ │
    │ Level 2: Manifest-Level Filtering (Manifest 级过滤) │
    │ ├── 利用 Manifest 中的 lower_bounds / upper_bounds │
    │ ├── 无需打开数据文件 │
    │ └── 跳过整个 Manifest 条目(可能对应数十个数据文件) │
    │ │
    │ Level 3: File-Level Filtering (文件级过滤) │
    │ ├── 利用 Manifest 中每个 Data File 的列统计 │
    │ ├── min/max/null_count 精确到每个文件 │
    │ └── 跳过文件内不包含匹配数据的 Parquet 文件 │
    │ │
    │ (额外) Row Group 过滤: │
    │ ├── Parquet 内部的 Row Group 也有 min/max 统计 │
    │ └── 读取时进一步跳过不相关的 Row Group │
    │ │
    └─────────────────────────────────────────────────────────────┘

    7.5 缓存策略

    # Spark 配置优化读取性能
    spark.conf.set("spark.sql.catalog.demo.cache-enabled", "true")
    spark.conf.set("spark.sql.catalog.demo.cache.expiration-interval-ms", "60000")

    # 读取时指定并行度
    spark.sql("""
    SELECT /*+ COALESCE(10) */ *
    FROM demo.analytics.orders
    WHERE order_date >= '2024-01-01'
    """
    )

    # 使用 read.split.target-size 控制每个任务读取的数据量
    ALTER TABLE demo.analytics.orders SET TBLPROPERTIES (
    'read.split.target-size' = '268435456' 256MB
    );

    ⚠️ 生产踩坑提醒:Time Travel 查询的是历史快照,不是当前数据。如果你的表有频繁的过期清理操作,历史快照可能已被删除,Time Travel 查询会失败。建议至少保留 7 天的快照。


    8. PyIceberg Python SDK 入门

    8.1 安装与配置

    # 安装 PyIceberg(0.11.0 是当前最新版本)
    pip install "pyiceberg[s3,glue]"

    # 或者安装完整版
    pip install "pyiceberg[s3,glue,rest,dynamodb,sql-postgres]"

    8.2 连接 Catalog

    from pyiceberg.catalog import load_catalog

    # 连接 REST Catalog
    catalog = load_catalog(
    "demo",
    **{
    "type": "rest",
    "uri": "http://localhost:8181",
    "warehouse": "s3://warehouse",
    "s3.endpoint": "http://localhost:9000",
    "s3.access-key-id": "admin",
    "s3.secret-access-key": "localpassword",
    "s3.path-style-access": "true",
    }
    )

    # 连接 AWS Glue Catalog
    catalog = load_catalog(
    "glue_catalog",
    **{
    "type": "glue",
    "warehouse": "s3://my-bucket/warehouse",
    "glue.region": "us-east-1",
    }
    )

    # 连接 Hive Catalog
    catalog = load_catalog(
    "hive_catalog",
    **{
    "type": "hive",
    "uri": "thrift://localhost:9083",
    "warehouse": "s3://warehouse",
    "s3.endpoint": "http://localhost:9000",
    }
    )

    8.3 命名空间与表管理

    from pyiceberg.schema import Schema
    from pyiceberg.types import (
    NestedField, StringType, LongType, TimestampType,
    DoubleType, DateType, BooleanType
    )
    from pyiceberg.partitioning import PartitionSpec, PartitionField
    from pyiceberg.transforms import DayTransform, BucketTransform, HourTransform
    from pyiceberg.sorting import SortOrder, SortField

    # 创建命名空间
    catalog.create_namespace("analytics")

    # 列出所有命名空间
    namespaces = catalog.list_namespaces()
    print(f"命名空间: {namespaces}")

    # 创建 Schema
    schema = Schema(
    NestedField(1, "order_id", LongType(), required=True),
    NestedField(2, "customer_id", LongType(), required=True),
    NestedField(3, "product_name", StringType(), required=False),
    NestedField(4, "amount", DoubleType(), required=False),
    NestedField(5, "order_date", DateType(), required=True),
    NestedField(6, "created_at", TimestampType(), required=False),
    NestedField(7, "is_active", BooleanType(), required=False),
    )

    # 定义分区规范
    partition_spec = PartitionSpec(
    PartitionField(5, 1000, DayTransform(), "order_date_day")
    )

    # 定义排序顺序
    sort_order = SortOrder(
    SortField(2, transform=BucketTransform(64)), # 按 customer_id 分桶排序
    SortField(5), # 再按 order_date 排序
    )

    # 创建表
    table = catalog.create_table(
    identifier="analytics.orders",
    schema=schema,
    partition_spec=partition_spec,
    sort_order=sort_order,
    properties={
    "format-version": "3",
    "write.parquet.compression-codec": "zstd",
    }
    )

    # 列出所有表
    tables = catalog.list_tables("analytics")
    print(f"表列表: {tables}")

    # 加载已有表
    table = catalog.load_table("analytics.orders")

    8.4 数据读写操作

    import pyarrow as pa
    from datetime import datetime, date

    # —- 使用 PyArrow Table 写入 —-
    data = pa.table({
    "order_id": [1, 2, 3, 4, 5],
    "customer_id": [100, 200, 100, 300, 200],
    "product_name": ["Widget", "Gadget", "Widget", "Doohickey", "Gadget"],
    "amount": [29.99, 49.99, 29.99, 9.99, 49.99],
    "order_date": [date(2024, 1, 15)] * 5,
    "created_at": [datetime.now()] * 5,
    "is_active": [True, True, False, True, True],
    })

    # 追加写入
    table.append(data)

    # 覆盖写入
    table.overwrite(data)

    # —- 读取数据 —-
    # 读取为 PyArrow Table
    arrow_table = table.scan().to_arrow()
    print(f"行数: {len(arrow_table)}")

    # 带过滤条件的读取
    filtered = table.scan(
    row_filter="customer_id = 100 AND amount > 20",
    selected_fields=("order_id", "product_name", "amount"),
    ).to_arrow()

    # 读取为 Pandas DataFrame
    pandas_df = table.scan(
    row_filter="order_date >= '2024-01-15'"
    ).to_pandas()

    # 读取为 PyArrow RecordBatchReader(适合大数据集)
    reader = table.scan().to_arrow_batch_reader()
    for batch in reader:
    print(f"批次行数: {batch.num_rows}")

    8.5 Schema 操作

    from pyiceberg.table.update.schema import UpdateSchema

    # 查看当前 Schema
    print(table.schema())

    # 添加列
    with table.update_schema() as update:
    update.add_column("shipping_address", StringType())
    update.add_column("discount_rate", DoubleType())

    # 重命名列
    with table.update_schema() as update:
    update.rename_column("shipping_address", "delivery_address")

    # 删除列
    with table.update_schema() as update:
    update.delete_column("discount_rate")

    # 修改列类型(拓宽)
    with table.update_schema() as update:
    update.update_column("amount", new_type=DoubleType())

    # 查看 Schema 历史
    print(f"当前 Schema ID: {table.schema_id}")
    for schema_id, schema in table.schemas.items():
    print(f" Schema {schema_id}: {schema}")

    8.6 表属性与快照管理

    # 查看表属性
    print(table.properties)

    # 更新表属性
    with table.update_properties() as update:
    update.set("write.parquet.compression-codec", "zstd")
    update.set("history.expire.max-snapshot-age-ms", "604800000") # 7天

    # 查看快照
    for snapshot in table.snapshots:
    print(f"Snapshot {snapshot.snapshot_id}: "
    f"operation={snapshot.summary.operation}, "
    f"timestamp={snapshot.timestamp_ms}")

    # 查看当前快照
    print(f"当前快照: {table.current_snapshot().snapshot_id}")

    # 查看快照详情
    print(f"Manifest List: {table.current_snapshot().manifest_list}")

    🏆 面试加分提示:PyIceberg 0.11.0 是一个重大版本,包含 380+ PR,50+ 贡献者。它支持完整的 v3 读写操作,包括 Deletion Vectors 和 VARIANT 类型。PyIceberg 不依赖 JVM,是纯 Python 实现,非常适合数据工程和 ML 场景。


    第二部分:进阶篇


    9. Iceberg v3 新特性深度解读

    9.1 版本总览

    Iceberg v3 是自 v2(2022 年)以来最重大的规范更新,于 2025 年末批准,2026 年各引擎陆续全面支持:

    特性描述影响
    Deletion Vectors Roaring Bitmap 替代位置删除文件 DML 性能提升 10 倍
    Row Lineage _row_id + _last_updated_sequence_number 原生 CDC 支持
    VARIANT 类型 原生半结构化数据列 替代 JSON STRING
    默认列值 Schema 元数据记录默认值 消除数据回填
    纳秒时间戳 timestamp_ns / timestamptz_ns 高精度时间场景
    Geometry/Geography 空间数据类型 GIS 场景原生支持
    多参数分区转换 复合列分桶 更灵活的分区策略

    9.2 Deletion Vectors 深度解析

    背景:v2 位置删除的问题

    v2 位置删除文件(Position Delete Files)的问题:

    原始数据文件: file_a.parquet (10000 行)

    第 1 次删除 → delete_file_1.avro:
    file_path | pos
    file_a | 5
    file_a | 42
    file_a | 100

    第 2 次删除 → delete_file_2.avro:
    file_path | pos
    file_a | 200
    file_a | 350

    第 N 次删除 → delete_file_N.avro:

    读取时需要合并所有删除文件 → 性能随删除次数线性下降!

    v3 解决方案:Deletion Vectors

    v3 Deletion Vector(Roaring Bitmap):

    原始数据文件: file_a.parquet (10000 行)

    Deletion Vector(存储在 Puffin 文件中):
    ┌─────────────────────────────────────────────┐
    │ file_a.parquet → deletion_vector │
    │ bitmap: {5, 42, 100, 200, 350} │
    │ 压缩后大小: ~100 bytes(而非 N 个文件) │
    └─────────────────────────────────────────────┘

    读取时:
    1. 打开 data file
    2. 加载 deletion vector bitmap
    3. 对于每行 row_idx: if bitmap.contains(row_idx) → skip
    4. 位图检查是 O(1),性能恒定!

    — v3 表自动使用 Deletion Vectors
    CREATE TABLE demo.analytics.cdc_orders (
    order_id BIGINT,
    customer_id BIGINT,
    status STRING,
    updated_at TIMESTAMP
    ) USING iceberg
    TBLPROPERTIES (
    'format-version' = '3',
    'write.delete.mode' = 'merge-on-read',
    'write.update.mode' = 'merge-on-read'
    );

    — 删除操作会自动创建 Deletion Vector 而非位置删除文件
    DELETE FROM demo.analytics.cdc_orders WHERE order_id = 1001;

    Puffin 文件格式:

    Puffin 文件结构(存储 Deletion Vectors):

    ┌──────────────────────────────────────────────┐
    │ Magic: PUFFIN │
    ├──────────────────────────────────────────────┤
    │ Blob 1: deletion-vector for file_a.parquet │
    │ ├── type: "deletion-vector-v1" │
    │ ├── input-fields: [file_path] │
    │ └── data: Roaring Bitmap bytes │
    ├──────────────────────────────────────────────┤
    │ Blob 2: deletion-vector for file_b.parquet │
    │ └── … │
    ├──────────────────────────────────────────────┤
    │ Footer: │
    │ ├── blob metadata (offset, length, type) │
    │ └── footer length │
    ├──────────────────────────────────────────────┤
    │ Footer Length (4 bytes) │
    ├──────────────────────────────────────────────┤
    │ Magic: PUFFIN │
    └──────────────────────────────────────────────┘

    9.3 Row Lineage(行血缘)

    v3 为每一行分配两个不可见的元数据字段:

    字段类型说明
    _row_id BIGINT 行的唯一标识符,永久不变
    _last_updated_sequence_number BIGINT 最后修改该行的提交序列号

    Row Lineage 工作机制:

    初始写入 (sequence_number = 1):
    ┌────────┬──────────┬────────────┐
    │ _row_id│ name │ _last_seq │
    ├────────┼──────────┼────────────┤
    │ 1 │ Alice │ 1 │
    │ 2 │ Bob │ 1 │
    │ 3 │ Charlie │ 1 │
    └────────┴──────────┴────────────┘

    更新 Bob 的 name (sequence_number = 5):
    ┌────────┬──────────┬────────────┐
    │ _row_id│ name │ _last_seq │
    ├────────┼──────────┼────────────┤
    │ 1 │ Alice │ 1 │ ← 不变
    │ 2 │ Bobby │ 5 │ ← 更新!row_id 不变
    │ 3 │ Charlie │ 1 │ ← 不变
    └────────┴──────────┴────────────┘

    Compaction 后(物理文件重组,逻辑行不变):
    ┌────────┬──────────┬────────────┐
    │ _row_id│ name │ _last_seq │
    ├────────┼──────────┼────────────┤
    │ 1 │ Alice │ 1 │ ← row_id 保持
    │ 2 │ Bobby │ 5 │ ← row_id 保持
    │ 3 │ Charlie │ 1 │ ← row_id 保持
    └────────┴──────────┴────────────┘

    Row ID 分配机制(继承模型):

    并发写入时的 Row ID 分配:

    Table metadata: next_row_id = 1000

    Writer 1 (Flink subtask 0):
    └── file_1.parquet: first_row_id = 1000
    row_ids: null, null, null (写入时)
    → 读取时分配: 1000, 1001, 1002

    Writer 2 (Flink subtask 1):
    └── file_2.parquet: first_row_id = 1003
    row_ids: null, null (写入时)
    → 读取时分配: 1003, 1004

    无需同步协调,每个文件独立分配 row_id 范围!

    9.4 VARIANT 类型

    VARIANT 是 Iceberg v3 引入的原生半结构化数据类型,用于替代 JSON STRING:

    — 创建包含 VARIANT 列的表
    CREATE TABLE demo.analytics.raw_events (
    event_id BIGINT,
    event_type STRING,
    payload VARIANT, — 原生 VARIANT 类型
    event_time TIMESTAMP
    ) USING iceberg
    TBLPROPERTIES ('format-version' = '3');

    — 插入 VARIANT 数据
    INSERT INTO demo.analytics.raw_events VALUES
    (1, 'user_action',
    parse_json('{"user_id": 1001, "action": "click", "page": "/home", "metadata": {"browser": "Chrome"}}'),
    current_timestamp());

    — 查询 VARIANT 中的字段(列式性能!)
    SELECT
    event_id,
    payload:user_id AS user_id,
    payload:action AS action,
    payload:metadata:browser AS browser
    FROM demo.analytics.raw_events
    WHERE payload:action = 'click';

    VARIANT vs JSON STRING 对比:

    维度VARIANT 类型JSON STRING
    存储格式 二进制编码(紧凑) 文本(UTF-8)
    查询性能 列式性能(Shredding) 每次扫描需解析
    Filter Pushdown ✅ 支持 ❌ 不支持
    类型推断 自动 需要手动 CAST
    Schema 灵活性 每行可以不同结构 每行可以不同结构
    存储效率 高(Shredding 优化) 低(重复键名)
    Parquet 支持 ✅ 原生 ✅ 但无优化

    Shredding 优化:

    VARIANT Shredding( shredding = 碎片化存储):

    原始 VARIANT 值:
    {"user_id": 1001, "action": "click", "metadata": {"browser": "Chrome"}}

    Shredding 后的存储:
    ┌─────────────────────────────────────────────────┐
    │ variant_typed_value: (二进制编码) │
    │ ├── typed_field: "user_id" → 1001 (bigint) │
    │ ├── typed_field: "action" → "click" (string) │
    │ └── typed_field: "metadata" → (nested variant)│
    │ ├── typed_field: "browser" → "Chrome" │
    │ └── … │
    └─────────────────────────────────────────────────┘

    频繁查询的字段会被 Shredding 为独立的类型化列
    → 查询这些字段时享受列式存储的性能
    → 不查询的字段仍保留在 VARIANT 二进制中

    9.5 默认列值

    v3 允许在 Schema 中定义列的默认值,对历史行透明应用:

    — 添加带默认值的列(历史行自动获得默认值)
    ALTER TABLE demo.analytics.orders
    ADD COLUMNS (region STRING DEFAULT 'unknown');

    — 历史行查询时透明应用默认值
    — 无需数据回填!

    — 修改已有列的默认值
    ALTER TABLE demo.analytics.orders
    ALTER COLUMN status SET DEFAULT 'pending';

    — 添加 NOT NULL 列(有默认值时可以做到!)
    ALTER TABLE demo.analytics.orders
    ADD COLUMNS (is_active BOOLEAN NOT NULL DEFAULT true);

    默认列值工作原理:

    Metadata 中记录:
    schema.fields[region].initial-default = "unknown"
    schema.fields[region].write-default = "unknown"

    读取时:
    ┌──────────────────────────────────────┐
    │ Data File (旧版本,无 region 列) │
    │ id | name | amount │
    │ 1 | 甲 | 100 │
    │ 2 | 乙 | 200 │
    └──────────────────────────────────────┘
    ↓ 读取时透明添加默认值
    ┌──────────────────────────────────────┐
    │ id | name | amount | region │
    │ 1 | 甲 | 100 | "unknown" │ ← 自动填充
    │ 2 | 乙 | 200 | "unknown" │ ← 自动填充
    └──────────────────────────────────────┘

    9.6 纳秒时间戳与空间类型

    — 纳秒时间戳(适用于高频交易、IoT 传感器等场景)
    CREATE TABLE demo.iot.sensor_data (
    sensor_id BIGINT,
    reading_time_ns TIMESTAMP_NTZ_NS, — 纳秒精度,无时区
    temperature DOUBLE,
    location GEOMETRY — 空间数据类型
    ) USING iceberg
    TBLPROPERTIES ('format-version' = '3');

    — 带时区的纳秒时间戳
    CREATE TABLE demo.finance.trades (
    trade_id BIGINT,
    trade_time TIMESTAMP_TZ_NS, — 纳秒精度,含时区
    amount DECIMAL(20, 8),
    counterparty_id STRING
    ) USING iceberg
    TBLPROPERTIES ('format-version' = '3');

    — 空间数据查询示例
    SELECT sensor_id, temperature
    FROM demo.iot.sensor_data
    WHERE ST_Within(location, ST_GeomFromText('POLYGON((…))'));

    9.7 从 v2 升级到 v3

    — 升级表格式版本(元数据操作,无需重写数据)
    ALTER TABLE demo.analytics.orders SET TBLPROPERTIES (
    'format-version' = '3'
    );

    — 验证升级成功
    SELECT * FROM demo.analytics.orders.properties;

    — 注意事项:
    — 1. 升级是元数据操作,不重写数据文件
    — 2. v2 的数据文件在 v3 中仍然有效
    — 3. 新写入的文件使用 v3 特性
    — 4. 确保所有读写引擎支持 v3!

    引擎v3 支持状态(2026 年中)
    Spark 4.0 + Iceberg 1.11.0 ✅ 完整读写支持
    Flink 2.0+ + Iceberg 1.10+ ✅ Deletion Vectors + Row Lineage
    Snowflake ✅ GA(2026年5月7日)
    Databricks Runtime 18.0+ ✅ 完整支持
    Trino ✅ 读写 DV,Row Lineage 部分支持
    StarRocks ✅ 完整兼容 v3
    Apache Doris 4.1+ ✅ UPDATE/DELETE/MERGE + v3

    ⚠️ 生产踩坑提醒:v3 升级路径是非破坏性的,现有 v2 数据文件在 v3 表中仍然有效。但如果有任何引擎不支持 v3 写入(比如旧版 Trino),不要升级该表,保持 v2 直到所有引擎就绪。VARIANT 类型在 Parquet 中才有 Shredding 支持,Avro 和 ORC 中只支持基本二进制编码。


    10. Iceberg 与计算引擎集成

    10.1 Spark 4.0 + Iceberg

    Spark 4.0 搭配 Iceberg 1.11.0 是当前最完整的开源 Iceberg 实现:

    from pyspark.sql import SparkSession

    spark = SparkSession.builder \\
    .appName("Spark4Iceberg") \\
    .config("spark.sql.extensions",
    "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \\
    .config("spark.sql.catalog.lake", "org.apache.iceberg.spark.SparkCatalog") \\
    .config("spark.sql.catalog.lake.type", "rest") \\
    .config("spark.sql.catalog.lake.uri", "http://polaris:8181") \\
    .config("spark.sql.catalog.lake.warehouse", "s3://lakehouse") \\
    .config("spark.sql.catalog.lake.io-impl",
    "org.apache.iceberg.aws.s3.S3FileIO") \\
    # Spark 4.0 优化
    .config("spark.sql.iceberg.planning.mode", "distributed") \\
    .config("spark.sql.catalog.lake.head-tracking.enabled", "true") \\
    .getOrCreate()

    # 批量写入优化
    spark.conf.set("spark.sql.catalog.lake.write.target-file-size-bytes", "536870912")
    spark.conf.set("spark.sql.catalog.lake.write.distribution-mode", "hash")

    # 读取优化
    spark.conf.set("spark.sql.catalog.lake.read.split.target-size", "134217728")

    Spark + Iceberg 关键配置:

    配置项推荐值说明
    write.distribution-mode hash 写入分布模式(none/hash/range)
    write.target-file-size-bytes 536870912 (512MB) 目标文件大小
    read.split.target-size 134217728 (128MB) 读取分片目标大小
    planning.mode distributed 文件规划模式(local/distributed)

    10.2 Flink + Iceberg(流式写入/读取)

    — Flink SQL Client 中注册 Iceberg Catalog
    CREATE CATALOG lake WITH (
    'type' = 'iceberg',
    'catalog-type' = 'rest',
    'uri' = 'http://polaris:8181',
    'warehouse' = 's3://lakehouse',
    'io-impl' = 'org.apache.iceberg.aws.s3.S3FileIO'
    );

    — 使用 Catalog
    USE CATALOG lake;
    CREATE DATABASE IF NOT EXISTS streaming;

    — 创建 Flink 写入表(启用 checkpoint)
    CREATE TABLE streaming.kafka_events (
    event_id BIGINT,
    event_type STRING,
    payload STRING,
    event_time TIMESTAMP(3),
    WATERMARK FOR event_time AS event_time INTERVAL '5' SECOND
    ) WITH (
    'connector' = 'kafka',
    'topic' = 'events',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json',
    'scan.startup.mode' = 'earliest-offset');

    — 将 Kafka 数据实时写入 Iceberg
    CREATE TABLE streaming.iceberg_events (
    event_id BIGINT,
    event_type STRING,
    payload STRING,
    event_time TIMESTAMP(3)
    );

    — Flink 实时写入
    INSERT INTO streaming.iceberg_events
    SELECT event_id, event_type, payload, event_time
    FROM streaming.kafka_events;

    Flink + Iceberg 关键配置:

    # flink-conf.yaml 关键配置
    execution.checkpointing.interval: 30s # checkpoint 间隔 = 数据新鲜度
    execution.checkpointing.min-pause: 10s
    state.backend: rocksdb
    state.checkpoints.dir: s3://lakehouse/checkpoints

    # Iceberg 写入优化
    table.exec.sink.upsert-materialize: none
    table.exec.sink.not-null-enforcer: drop

    Flink + Iceberg 流式写入流程:

    Kafka Topic Flink Job Iceberg Table
    ┌─────────┐ Checkpoint ┌─────────┐ Commit ┌─────────┐
    │ offset 0│ ──────────→ │ Process │ ──────────→ │Snapshot │
    │ offset 1│ 每 30s 一次 │ Buffer │ 新快照 │ N+1 │
    │ offset 2│ │ Write │ │ │
    │ … │ │ Commit │ │ │
    └─────────┘ └─────────┘ └─────────┘

    Checkpoint 间隔 = 数据新鲜度 = 文件产生频率
    30s checkpoint × 16 parallelism = 32 files/min

    10.3 Trino + Iceberg

    — Trino 中配置 Iceberg Catalog
    — etc/catalog/lake.properties
    — connector.name=iceberg
    — iceberg.catalog.type=rest
    — iceberg.rest-catalog.uri=http://polaris:8181
    — iceberg.rest-catalog.warehouse=s3://lakehouse
    — iceberg.rest-catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO
    — fs.native-s3.enabled=true
    — s3.endpoint=http://minio:9000
    — s3.path-style-access=true
    — s3.aws-access-key=admin
    — s3.aws-secret-key=password

    — Trino SQL 查询 Iceberg 表
    SELECT
    customer_id,
    COUNT(*) AS order_count,
    SUM(total_amount) AS total_spent
    FROM lake.analytics.orders
    WHERE order_date >= DATE '2024-01-01'
    GROUP BY customer_id
    ORDER BY total_spent DESC
    LIMIT 100;

    — Trino 中的 Time Travel
    SELECT * FROM lake.analytics.orders
    FOR TIMESTAMP AS OF TIMESTAMP '2024-01-15 12:00:00';

    — Trino 中的分支查询
    SELECT * FROM lake.analytics.orders
    FOR VERSION AS OF 'experiment_branch';

    10.4 StarRocks/Doris + Iceberg(外表查询)

    — StarRocks External Catalog
    CREATE EXTERNAL CATALOG iceberg_catalog
    PROPERTIES (
    "type" = "iceberg",
    "iceberg.catalog.type" = "rest",
    "iceberg.catalog.uri" = "http://polaris:8181",
    "iceberg.catalog.warehouse" = "s3://lakehouse",
    "aws.s3.access_key" = "admin",
    "aws.s3.secret_key" = "password",
    "aws.s3.endpoint" = "http://minio:9000"
    );

    — 查询 Iceberg 表
    SELECT * FROM iceberg_catalog.analytics.orders
    WHERE order_date >= '2024-01-01'
    LIMIT 100;

    — Apache Doris 4.1+ 支持对 Iceberg 表执行 UPDATE/DELETE/MERGE
    — Doris 外部表 + Iceberg v3
    SELECT * FROM iceberg_catalog.analytics.orders
    WHERE _change_type IS NOT NULL;

    10.5 Snowflake + Iceberg

    Snowflake 于 2026 年 5 月 7 日正式 GA 了 Iceberg v3 支持:

    — Snowflake 中创建 Iceberg Table
    CREATE ICEBERG TABLE snowflake_db.orders (
    order_id BIGINT,
    customer_id BIGINT,
    total_amount DECIMAL(10, 2),
    order_date DATE,
    payload VARIANT, — Snowflake 原生 VARIANT
    location GEOGRAPHY — 空间类型
    )
    CATALOG_TABLE_NAME = 'analytics.orders'
    CATALOG = 'polaris_catalog'
    EXTERNAL_VOLUME = 's3_volume';

    — 直接查询
    SELECT * FROM snowflake_db.orders
    WHERE order_date >= '2024-01-01';

    — Snowflake 的 Managed Storage(2026 GA)
    — Snowflake 自动管理 Iceberg 表的 compaction、文件优化

    引擎最佳场景v3 完整度
    Spark 4.0 批处理、ETL、数据工程 ✅ 最完整
    Flink 2.0+ 流式写入、CDC 入湖 ✅ DV + Row Lineage
    Trino 交互式查询、BI ✅ DV 读写,RL 部分
    Snowflake 云数仓、数据共享 ✅ 完整 v3 GA
    Databricks 18.0+ Delta 生态 + Iceberg 兼容 ✅ 完整
    StarRocks OLAP 加速查询 ✅ 外表查询
    Doris 4.1+ 单系统查改维 ✅ 含 UPDATE/MERGE

    11. Catalog 选型与部署

    11.1 Catalog 对比

    特性Hive CatalogREST CatalogJDBC CatalogNessie CatalogAWS Glue
    元数据存储 Hive Metastore 独立服务 + DB 关系数据库 Nessie Server AWS Glue
    协议 Thrift HTTP/REST JDBC HTTP/REST AWS API
    ACID 提交 ⚠️ 有限 ✅ 服务端 ✅ 数据库事务 ✅ Git-like ✅ DynamoDB
    多表事务
    凭证代理 ✅ (STS)
    RBAC ⚠️ Ranger ✅ OAuth2 ✅ IAM
    Git 分支语义 ✅ 原生
    高可用 ⚠️ 需 HA 配置 ✅ 多实例 ✅ 依赖 DB ✅ 多实例 ✅ AWS 保证
    推荐度 (2026) ⭐⭐ 遗留 ⭐⭐⭐⭐⭐ 首选 ⭐⭐⭐ ⭐⭐⭐ ⭐⭐⭐⭐ AWS 生态
    开源项目 Hive Polaris / Lakekeeper Project Nessie AWS 托管

    11.2 REST Catalog 架构

    ┌─────────────────────────────────────────────────────────────────┐
    │ REST Catalog 架构 │
    ├─────────────────────────────────────────────────────────────────┤
    │ │
    │ ┌─────────┐ ┌──────────┐ ┌────────┐ ┌──────────┐ │
    │ │ Spark │ │ Flink │ │ Trino │ │PyIceberg │ │
    │ └────┬────┘ └────┬─────┘ └───┬────┘ └────┬─────┘ │
    │ │ │ │ │ │
    │ └────────────┴─────┬──────┴─────────────┘ │
    │ │ REST API (HTTP/JSON) │
    │ │ │
    │ ┌───────────┴───────────┐ │
    │ │ REST Catalog 服务 │ │
    │ │ (Polaris/Lakekeeper) │ │
    │ │ │ │
    │ │ ┌─────────────────┐ │ │
    │ │ │ 认证 (OAuth2) │ │ │
    │ │ │ 凭证代理 (STS) │ │ │
    │ │ │ 服务端提交 │ │ │
    │ │ │ RBAC 权限控制 │ │ │
    │ │ └────────┬────────┘ │ │
    │ └───────────┼───────────┘ │
    │ │ │
    │ ┌───────────┴───────────┐ │
    │ │ 后端存储 │ │
    │ │ PostgreSQL / JDBC │ │
    │ └───────────────────────┘ │
    │ │
    └─────────────────────────────────────────────────────────────────┘

    11.3 Apache Polaris 部署

    Apache Polaris 于 2026 年 2 月毕业为 Apache 顶级项目,是推荐的 REST Catalog 实现:

    # docker-compose.yml for Polaris
    services:
    postgres:
    image: postgres:17
    environment:
    POSTGRES_USER: polaris
    POSTGRES_PASSWORD: polaris
    POSTGRES_DB: polaris
    volumes:
    pgdata:/var/lib/postgresql/data

    polaris:
    image: apache/polaris:latest
    depends_on:
    postgres
    ports:
    "8181:8181"
    environment:
    POLARIS_DB_URL: jdbc:postgresql://postgres:5432/polaris
    POLARIS_DB_USER: polaris
    POLARIS_DB_PASSWORD: polaris
    POLARIS_SERVICE_PORT: 8181

    volumes:
    pgdata:

    11.4 多 Catalog 联邦查询

    # PyIceberg 支持同时连接多个 Catalog
    from pyiceberg.catalog import load_catalog

    # 生产 Catalog
    prod_catalog = load_catalog("prod", type="rest", uri="http://prod-polaris:8181")

    # 开发 Catalog
    dev_catalog = load_catalog("dev", type="rest", uri="http://dev-polaris:8181")

    # 跨 Catalog 操作:从生产复制到开发
    prod_table = prod_catalog.load_table("analytics.orders")
    dev_table = dev_catalog.load_table("staging.orders_copy")

    # 读取生产数据,写入开发环境
    data = prod_table.scan(row_filter="order_date >= '2024-01-01'").to_arrow()
    dev_table.append(data)

    🏆 面试加分提示:2026 年 REST Catalog 已成为事实标准。Apache Polaris 是 Snowflake 和 Dremio 联合开发并贡献给 Apache 的开源项目。它支持凭证代理(Credential Vending)——引擎不需要直接访问存储的密钥,Catalog 为每次查询签发临时凭证(STS Token),实现零信任安全架构。

    ⚠️ 生产踩坑提醒:选择 Catalog 时考虑三个关键因素:(1) 你的计算引擎生态——Spark 生态优先 Polaris,AWS 生态优先 Glue;(2) 安全需求——需要凭证代理就必须用 REST Catalog;(3) 运维能力——Hive Catalog 需要运维 Thrift Server,REST Catalog 需要运维服务+数据库。


    12. 性能调优

    12.1 小文件合并(Compaction)

    小文件是 Iceberg 性能的头号杀手。每次 Flink checkpoint 或 Spark 微批次都会产生新文件:

    小文件问题示意:

    Flink 30s checkpoint × 16 parallelism:
    → 每分钟 32 个文件
    → 每小时 1,920 个文件
    → 每天 46,080 个文件!

    读取时需要打开数万个文件 → 查询耗时从秒级退化到分钟级

    Spark Actions 合并方案:

    from pyiceberg.catalog import load_catalog

    catalog = load_catalog("demo")
    table = catalog.load_table("analytics.events")

    # 重写数据文件(合并小文件)
    from pyspark.sql import SparkSession

    spark = SparkSession.builder.getOrCreate()

    # 使用 Spark Actions 进行 compaction
    spark.sql("""
    CALL demo.system.rewrite_data_files(
    table => 'analytics.events',
    options => map(
    'target-file-size-bytes', '536870912',
    'max-concurrent-file-group-rewrites', '5'
    )
    )
    """
    )

    # 使用 PyIceberg 进行 compaction(无需 Spark)
    # PyIceberg 0.11.0 支持基本的文件重写
    from pyiceberg.table import Table

    table.rewrite_data_files(
    target_size_in_bytes=536870912,
    max_concurrent_file_group_rewrites=5
    )

    Compaction 策略选择:

    策略说明适用场景
    BinPack 将小文件打包为大文件 通用场景
    Sort 合并时按指定列排序 需要排序优化查询
    Z-Order 多维 Z-Order 排序 多列组合过滤

    — 按特定列排序后重写(提升过滤查询性能)
    CALL demo.system.rewrite_data_files(
    table => 'analytics.events',
    strategy => 'sort',
    sort_order => 'event_time DESC NULLS LAST, user_id ASC',
    options => map('target-file-size-bytes', '536870912')
    );

    12.2 Snapshot 过期清理

    # 清理过期快照
    spark.sql("""
    CALL demo.system.expire_snapshots(
    table => 'analytics.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00',
    retain_last => 10
    )
    """
    )

    # 使用表属性自动过期
    ALTER TABLE demo.analytics.events SET TBLPROPERTIES (
    'history.expire.max-snapshot-age-ms' = '604800000', 7
    'history.expire.min-snapshots-to-keep' = '10' 至少保留10
    );

    12.3 Orphan Files 清理

    — 清理孤立文件(存在于存储但不在任何快照中引用的文件)
    CALL demo.system.remove_orphan_files(
    table => 'analytics.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00'
    );

    12.4 性能调优清单

    ┌─────────────────────────────────────────────────────────────┐
    │ Iceberg 性能调优清单 │
    ├─────────────────────────────────────────────────────────────┤
    │ │
    │ 写入优化: │
    │ ├── 目标文件大小: 128MB ~ 512MB │
    │ ├── 压缩算法: ZSTD(最佳压缩/速度比) │
    │ ├── 写入并行度: 根据分区数调整 │
    │ ├── 写入分布模式: hash(均匀分布到分区) │
    │ └── 批量提交: 减少提交频率,增大每批次数据量 │
    │ │
    │ 读取优化: │
    │ ├── 分区裁剪: 确保查询条件命中分区键 │
    │ ├── 排序优化: 数据按查询常用列排序 │
    │ ├── 读取分片大小: 128MB~256MB per task │
    │ ├── 文件规划模式: distributed(大表) │
    │ └── 缓存: 启用 manifest cache │
    │ │
    │ 维护优化: │
    │ ├── 定期 Compaction: 每天或每 N 次写入后 │
    │ ├── 快照过期: 保留 7~30 天 │
    │ ├── 孤立文件清理: 每周执行 │
    │ ├── Manifest 文件合并: 自动或手动触发 │
    │ └── 监控: 文件数、快照数、元数据大小 │
    │ │
    └─────────────────────────────────────────────────────────────┘

    12.5 数据倾斜处理

    # 问题:某个分区的数据量远大于其他分区
    # 解决:使用 bucket 分区 + 合理的桶数

    # 方案1:增加桶数
    ALTER TABLE demo.analytics.events
    ADD PARTITION FIELD bucket(256, user_id); # 256个桶

    # 方案2:使用 truncate 分区减少热点
    # 适合 URL 前缀相似的日志数据

    # 方案3:写入时调整分布
    spark.conf.set("spark.sql.catalog.demo.write.distribution-mode", "hash")

    # 方案4:自定义写入并行度
    df.write \\
    .format("iceberg") \\
    .option("write.distribution-mode", "hash") \\
    .option("spark.sql.files.maxRecordsPerFile", "1000000") \\
    .mode("append") \\
    .saveAsTable("demo.analytics.events")

    ⚠️ 生产踩坑提醒:Compaction 是最容易被忽略的维护操作。很多团队部署了 Iceberg 却从不 compaction,几个月后查询性能断崖式下降。建议:(1) 配置自动 compaction 调度(每天凌晨执行);(2) 监控每个表的文件数量告警阈值;(3) 流式写入场景尤其需要频繁 compaction。


    13. 数据湖 CDC 与流批一体

    13.1 CDC 数据入湖架构

    ┌─────────────────────────────────────────────────────────────────┐
    │ CDC 数据入湖完整架构 │
    ├─────────────────────────────────────────────────────────────────┤
    │ │
    │ 源数据库 传输层 计算层 │
    │ ┌──────────┐ ┌─────────┐ ┌──────────┐ │
    │ │ MySQL │ │ │ │ │ │
    │ │ 订单表 │──── CDC ───→│ Kafka │──────→│ Flink │ │
    │ │ 用户表 │ Debezium │ Topic │ │ 实时 │ │
    │ │ 商品表 │ │ │ │ 合并 │ │
    │ └──────────┘ └─────────┘ └────┬─────┘ │
    │ │ │
    │ MERGE INTO │
    │ │ │
    │ ┌─────┴─────┐ │
    │ │ Iceberg │ │
    │ │ 数据湖 │ │
    │ │ (v3 MoR) │ │
    │ └─────┬─────┘ │
    │ │ │
    │ ┌──────────┼─────────┐│
    │ │ │ ││
    │ ┌────┴──┐ ┌────┴──┐ ┌───┴┐│
    │ │ Trino │ │Spark │ │Snow││
    │ │ (OLAP)│ │(batch)│ │flk ││
    │ └───────┘ └───────┘ └────┘│
    │ │
    └─────────────────────────────────────────────────────────────────┘

    13.2 Debezium + Flink + Iceberg 实战

    Step 1: 配置 Debezium CDC 连接器

    — Flink SQL: 创建 Kafka Source(Debezium JSON 格式)
    CREATE TABLE mysql_orders_source (
    order_id BIGINT,
    customer_id BIGINT,
    order_date DATE,
    total_amount DECIMAL(10, 2),
    status STRING,
    created_at TIMESTAMP(3),
    updated_at TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
    ) WITH (
    'connector' = 'kafka',
    'topic' = 'mysql.orders.binlog',
    'properties.bootstrap.servers' = 'kafka:9092',
    'properties.group.id' = 'flink-iceberg-cdc',
    'scan.startup.mode' = 'earliest-offset',
    'format' = 'debezium-json',
    'debezium-json.schema-include' = 'true'
    );

    Step 2: 创建 Iceberg 目标表

    — 创建 Iceberg 目标表(v3 格式,merge-on-read)
    CREATE TABLE lake.cdc.orders (
    order_id BIGINT,
    customer_id BIGINT,
    order_date DATE,
    total_amount DECIMAL(10, 2),
    status STRING,
    created_at TIMESTAMP(3),
    updated_at TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
    ) WITH (
    'format-version' = '3',
    'write.delete.mode' = 'merge-on-read',
    'write.update.mode' = 'merge-on-read',
    'write.merge.mode' = 'merge-on-read',
    'write.parquet.compression-codec' = 'zstd'
    );

    Step 3: Flink 实时写入 Iceberg

    // Flink Java: 使用 FlinkSink API 写入 Iceberg
    import org.apache.iceberg.flink.sink.FlinkSink;
    import org.apache.iceberg.flink.TableLoader;

    FlinkSink.forRowData(inputDataStream)
    .tableLoader(TableLoader.fromCatalog(catalog, tableIdentifier))
    .overwrite(false)
    .distributionMode(DistributionMode.HASH)
    .writeParallelism(16)
    .build();

    13.3 v3 Row Lineage 与原生 CDC

    v3 的 Row Lineage 特性为 CDC 带来了原生支持:

    # 利用 Row Lineage 进行增量处理
    from pyiceberg.catalog import load_catalog

    catalog = load_catalog("demo")
    table = catalog.load_table("cdc.orders")

    # 获取上次处理的序列号
    last_sequence = get_last_processed_sequence() # 从外部存储获取

    # 查询自上次处理以来变更的行
    # _last_updated_sequence_number > last_sequence
    changed_rows = table.scan(
    row_filter=f"_last_updated_sequence_number > {last_sequence}"
    ).to_arrow()

    # 处理变更的行
    for row in changed_rows.to_pylist():
    row_id = row['_row_id']
    seq_num = row['_last_updated_sequence_number']
    # … 增量处理逻辑

    13.4 流批一体架构设计

    ┌─────────────────────────────────────────────────────────────────┐
    │ 流批一体架构 │
    ├─────────────────────────────────────────────────────────────────┤
    │ │
    │ 写入路径: │
    │ ┌──────────┐ ┌──────────┐ ┌──────────────────┐ │
    │ │ 实时流 │ ──→ │ Flink │ ──→ │ Iceberg Table │ │
    │ │ (Kafka) │ │ (流写入) │ │ (统一存储层) │ │
    │ └──────────┘ └──────────┘ └────────┬─────────┘ │
    │ │ │
    │ 读取路径: │ │
    │ ┌──────────┐ │ │
    │ │ 实时查询 │ ←──── 增量读取(新快照)────────┤ │
    │ │ (Flink) │ │ │
    │ └──────────┘ │ │
    │ ┌──────────┐ │ │
    │ │ 批处理 │ ←──── 全量读取(Time Travel)───┤ │
    │ │ (Spark) │ │ │
    │ └──────────┘ │ │
    │ ┌──────────┐ │ │
    │ │ 交互查询 │ ←──── OLAP 查询 ───────────────┤ │
    │ │ (Trino) │ │ │
    │ └──────────┘ │ │
    │ │
    │ 核心价值: │
    │ ├── 一份数据,多种访问模式 │
    │ ├── 实时和离线使用同一张表 │
    │ ├── 消除数据冗余和同步延迟 │
    │ └── ACID 保证并发安全 │
    │ │
    └─────────────────────────────────────────────────────────────────┘

    13.5 Exactly-Once 保证

    Flink + Iceberg 通过 Checkpoint 协调实现 Exactly-Once:

    Flink Checkpoint 与 Iceberg 快照的关系:

    Checkpoint 1 ──→ Iceberg Snapshot 1
    ├── offset 0~1000 的数据写入 data_file_1
    ├── data_file_1 记录在 manifest_1 中
    ├── manifest_1 记录在 manifest_list_1 中
    └── manifest_list_1 原子提交到 catalog

    Checkpoint 2 ──→ Iceberg Snapshot 2
    ├── offset 1001~2000 的数据写入 data_file_2
    └── … 原子提交

    如果 Checkpoint 2 失败:
    ├── offset 1001~2000 的数据被丢弃
    ├── Flink 从 Checkpoint 1 恢复
    └── 重新消费 offset 1001~2000 → 幂等写入

    关键:每个 Checkpoint 对应一个原子提交
    要么全部成功(新快照可见),要么全部回滚

    ⚠️ 生产踩坑提醒:CDC 入湖场景下,Checkpoint 间隔的设置至关重要。30 秒间隔意味着 30 秒数据新鲜度,但也意味着每分钟 32 个文件(16 parallelism)。建议:(1) 初始设置 60 秒间隔,观察文件增长;(2) 配合自动 compaction 控制文件数量;(3) v3 的 Deletion Vectors 可以显著降低 MOR 读取开销,推荐使用。


    14. 表维护与生命周期管理

    14.1 自动化维护方案

    # Spark 维护任务(推荐通过 Airflow/DolphinScheduler 调度)
    from pyspark.sql import SparkSession

    spark = SparkSession.builder \\
    .appName("IcebergMaintenance") \\
    .config("spark.sql.extensions",
    "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \\
    .config("spark.sql.catalog.demo", "org.apache.iceberg.spark.SparkCatalog") \\
    .config("spark.sql.catalog.demo.type", "rest") \\
    .config("spark.sql.catalog.demo.uri", "http://polaris:8181") \\
    .getOrCreate()

    # —- 1. 小文件合并(每天执行)—-
    spark.sql("""
    CALL demo.system.rewrite_data_files(
    table => 'analytics.events',
    options => map(
    'target-file-size-bytes', '536870912',
    'max-concurrent-file-group-rewrites', '5',
    'partial-progress.enabled', 'true'
    )
    )
    """
    )

    # —- 2. 过期快照清理(每天执行)—-
    spark.sql("""
    CALL demo.system.expire_snapshots(
    table => 'analytics.events',
    older_than => TIMESTAMP '2024-01-08 00:00:00',
    retain_last => 10,
    stream_results => true
    )
    """
    )

    # —- 3. 孤立文件清理(每周执行)—-
    spark.sql("""
    CALL demo.system.remove_orphan_files(
    table => 'analytics.events',
    older_than => TIMESTAMP '2024-01-01 00:00:00'
    )
    """
    )

    # —- 4. 删除文件合并(Rewrite Delete Files)—-
    # 将多个小的 Deletion Vector 合并为更紧凑的形式
    spark.sql("""
    CALL demo.system.rewrite_position_delete_files(
    table => 'analytics.events',
    options => map('max-concurrent-file-group-rewrites', '3')
    )
    """
    )

    14.2 维护调度策略

    维护任务频率执行时间优先级
    小文件合并 每天 1~4 次 凌晨低峰
    快照过期清理 每天 1 次 凌晨
    孤立文件清理 每周 1 次 周末
    删除文件合并 每天 1 次 凌晨 高(MOR 表)
    Manifest 合并 自动

    14.3 监控告警

    # PyIceberg 健康检查脚本
    from pyiceberg.catalog import load_catalog
    from datetime import datetime

    catalog = load_catalog("demo")
    table = catalog.load_table("analytics.events")

    # 收集表健康指标
    current_snapshot = table.current_snapshot()
    all_snapshots = list(table.snapshots)
    num_snapshots = len(all_snapshots)

    # 计算快照年龄
    now_ms = int(datetime.now().timestamp() * 1000)
    oldest_snapshot_age_days = (now_ms all_snapshots[0].timestamp_ms) / (1000 * 86400) if all_snapshots else 0
    newest_snapshot_age_days = (now_ms current_snapshot.timestamp_ms) / (1000 * 86400) if current_snapshot else 0

    # 获取 manifest 文件数量(从快照摘要)
    summary = current_snapshot.summary if current_snapshot else {}
    num_data_files = int(summary.get('total-data-files', 0))
    num_delete_files = int(summary.get('total-delete-files', 0))

    print(f"""
    === Iceberg 表健康报告 ===
    表: analytics.events
    快照数:
    {num_snapshots}
    最旧快照:
    {oldest_snapshot_age_days:.1f} 天前
    最新快照:
    {newest_snapshot_age_days:.1f} 小时前
    数据文件数:
    {num_data_files}
    删除文件数:
    {num_delete_files}
    """
    )

    # 告警条件
    if num_data_files > 10000:
    print("⚠️ 警告:数据文件数超过 10000,建议执行 compaction!")
    if num_snapshots > 1000:
    print("⚠️ 警告:快照数超过 1000,建议清理过期快照!")
    if num_delete_files > 1000:
    print("⚠️ 警告:删除文件数超过 1000,建议执行 delete file compaction!")

    14.4 存储分层(冷热数据)

    # 使用 Iceberg 表属性 + S3 生命周期规则实现冷热分层
    # 热数据(最近 7 天)→ S3 Standard
    # 温数据(7~90 天)→ S3 Standard-IA
    # 冷数据(>90 天)→ S3 Glacier

    # Iceberg 端:标记分区时间范围
    ALTER TABLE demo.analytics.events SET TBLPROPERTIES (
    'write.data.layout' = 'compact',
    'history.expire.max-snapshot-age-ms' = '7776000000' # 90天
    );

    # S3 端:配置生命周期规则(AWS Console 或 Terraform)
    # 规则示例:
    # – 前缀: warehouse/analytics/events/data/
    # – 30天后转为 Standard-IA
    # – 90天后转为 Glacier Instant Retrieval
    # – 365天后转为 Glacier Deep Archive

    ⚠️ 生产踩坑提醒:表维护的"黄金法则"是——写入路径和维护路径必须同时部署。最常见的失败是只部署了写入路径而没有维护路径。Freshness 如果不能被查询到,就没有意义。建议将维护任务纳入 CI/CD 流水线,与数据管道一起部署。


    15. 安全与权限

    15.1 安全架构总览

    ┌─────────────────────────────────────────────────────────────┐
    │ Iceberg 安全架构 │
    ├─────────────────────────────────────────────────────────────┤
    │ │
    │ 认证层 (Authentication): │
    │ ├── Kerberos (企业环境) │
    │ ├── OAuth2 / JWT (REST Catalog) │
    │ ├── IAM Role (AWS) │
    │ └── LDAP / SAML (企业 SSO) │
    │ │
    │ 授权层 (Authorization): │
    │ ├── Catalog 级别: 命名空间 / 表的 CRUD 权限 │
    │ ├── 表级别: SELECT / INSERT / DELETE / ALTER │
    │ ├── 列级别: 指定用户可见的列 │
    │ ├── 行级别: 基于条件的行过滤 │
    │ └── RBAC / ABAC (基于角色/属性的访问控制) │
    │ │
    │ 数据安全: │
    │ ├── 静态加密 (Encryption at Rest): S3 SSE / KMS │
    │ ├── 传输加密 (Encryption in Transit): TLS/SSL │
    │ ├── v3 加密密钥管理: 表级别加密密钥 │
    │ └── 数据脱敏: 动态脱敏 / 静态脱敏 │
    │ │
    └─────────────────────────────────────────────────────────────┘

    15.2 REST Catalog 鉴权(OAuth2 / JWT)

    from pyiceberg.catalog import load_catalog

    # OAuth2 认证连接 REST Catalog
    catalog = load_catalog(
    "secure_catalog",
    **{
    "type": "rest",
    "uri": "https://polaris.example.com",
    "warehouse": "s3://secure-warehouse",
    # OAuth2 配置
    "rest.auth.type": "oauth2",
    "rest.auth.oauth2.credential": "client_id",
    "rest.auth.oauth2.client-secret": "client_secret",
    "rest.auth.oauth2.token-url": "https://auth.example.com/token",
    "rest.auth.oauth2.scope": "read write",
    }
    )

    Apache Polaris 的三级权限模型:

    Polaris RBAC 权限模型:

    Principal (用户/服务账号)

    ├── PrincipalRole (角色)
    │ ├── "data_engineer"
    │ ├── "data_analyst"
    │ └── "data_admin"

    └── CatalogRole (目录角色)
    ├── "table_reader" → 对 analytics.orders 有 SELECT
    ├── "table_writer" → 对 analytics.orders 有 INSERT
    └── "table_admin" → 对 analytics.orders 有全部权限

    授权链:
    Principal → PrincipalRole → CatalogRole → 权限

    示例:
    user_alice → data_engineer → table_writer → INSERT on analytics.orders
    user_bob → data_analyst → table_reader → SELECT on analytics.orders

    15.3 凭证代理(Credential Vending)

    凭证代理工作流程:

    1. 引擎请求访问表
    Spark → REST Catalog: "我要读 analytics.orders"

    2. Catalog 验证身份 + 权限
    Catalog → OAuth2 Server: 验证 Token
    Catalog → Polaris RBAC: 检查权限

    3. Catalog 签发临时凭证
    Catalog → STS: AssumeRole(table-specific policy)
    STS → Catalog: 临时 AccessKey + SessionToken (15分钟有效)

    4. 引擎使用临时凭证访问存储
    Spark → S3: 使用临时凭证读取文件
    S3: 验证凭证 + Session Policy → 允许访问

    优势:
    ├── 引擎永远不持有永久密钥
    ├── 凭证自动过期(15分钟 TTL)
    ├── 权限范围限定到具体表前缀
    └── 审计追踪完整

    15.4 列级权限与行级过滤

    # 通过 REST Catalog + Trino 实现列级权限
    # Trino catalog 配置中设置视图

    # 列级权限:创建受限视图
    spark.sql("""
    CREATE VIEW demo.analytics.orders_public AS
    SELECT
    order_id,
    customer_id,
    total_amount,
    order_date,
    '***' AS shipping_address — 脱敏
    FROM demo.analytics.orders
    """
    )

    # 行级过滤:不同部门看到不同数据
    # 通过 Trino Row Filter 实现
    # trino/etc/catalog/lake.properties 中配置 row filter

    15.5 数据加密

    — v3 加密配置(表级别加密密钥)
    CREATE TABLE demo.analytics.sensitive_data (
    user_id BIGINT,
    ssn STRING,
    credit_card STRING
    ) USING iceberg
    TBLPROPERTIES (
    'format-version' = '3',
    'encryption.key-id' = 'arn:aws:kms:us-east-1:123456789:key/abc-123',
    'encryption.algorithm' = 'AES-256-GCM'
    );

    — S3 层面加密(推荐方式)
    — 在存储层启用 SSE-KMS,所有写入自动加密
    — Iceberg 元数据文件也受 S3 加密保护

    加密层级实现方式适用场景
    存储层加密 S3 SSE-S3 / SSE-KMS 所有数据(推荐)
    表级加密 v3 encryption keys 敏感数据表
    列级加密 应用层加密 + 密文存储 极高安全需求
    传输加密 TLS 1.3 所有网络通信

    ⚠️ 生产踩坑提醒:凭证代理是零信任架构的核心。没有凭证代理,每个引擎都需要直接持有存储密钥,一旦密钥泄露后果严重。2026 年 Apache Polaris 1.4 支持多云凭证代理——一个 Polaris 实例可以同时为 AWS S3、Azure ADLS、GCP GCS 签发凭证。


    16. 生产环境部署架构

    16.1 典型生产架构图

    ┌─────────────────────────────────────────────────────────────────────┐
    │ 典型生产环境架构 │
    ├─────────────────────────────────────────────────────────────────────┤
    │ │
    │ 数据源层: │
    │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
    │ │ MySQL │ │ Postgres │ │ Kafka │ │ S3 CSV │ │
    │ │ (CDC) │ │ (CDC) │ │ (Event) │ │ (Batch) │ │
    │ └────┬─────┘ └────┬─────┘ └────┬─────┘ └────┬─────┘ │
    │ │ │ │ │ │
    │ 入湖层: │
    │ ┌────┴────────────┴────┐ ┌────┴─────────────┴──┐ │
    │ │ Flink CDC Cluster │ │ Spark ETL Cluster │ │
    │ │ (实时入湖) │ │ (批量入湖) │ │
    │ └──────────┬───────────┘ └──────────┬──────────┘ │
    │ │ │ │
    │ ───────────┼──────────────────────────┼──────────────── │
    │ │ REST Catalog │ │
    │ │ ┌───────────────────────────┐ │
    │ │ │ Apache Polaris │ │
    │ │ │ (OAuth2 + 凭证代理) │ │
    │ │ │ Backend: PostgreSQL HA │ │
    │ │ └──────────────┬────────────┘ │
    │ │ │ │
    │ 存储层: │ │ │
    │ ┌──────────┴─────────────────┴──────────────────────┐ │
    │ │ Amazon S3 / MinIO │ │
    │ │ ┌─────────────────────────────────────────────┐ │ │
    │ │ │ warehouse/ │ │ │
    │ │ │ ├── bronze/ (原始数据) │ │ │
    │ │ │ ├── silver/ (清洗后数据) │ │ │
    │ │ │ └── gold/ (业务聚合数据) │ │ │
    │ │ └─────────────────────────────────────────────┘ │ │
    │ └───────────────────────────────────────────────────┘ │
    │ │ │
    │ 查询层: │ │
    │ ┌──────────┴──────────────────────────────────────────┐ │
    │ │ │ │
    │ │ ┌─────────┐ ┌──────────┐ ┌────────┐ ┌────────┐ │ │
    │ │ │ Trino │ │ Spark │ │Snowflake│ │StarRocks│ │ │
    │ │ │(OLAP) │ │ (ETL) │ │(共享) │ │(OLAP) │ │ │
    │ │ └────────┘ └──────────┘ └────────┘ └────────┘ │ │
    │ │ │ │
    │ └──────────────────────────────────────────────────────┘ │
    │ │
    │ 维护层: │
    │ ┌──────────────────────────────────────────────────────┐ │
    │ │ Airflow / DolphinScheduler │ │
    │ │ ├── 每日 compaction │ │
    │ │ ├── 每日 snapshot 过期清理 │ │
    │ │ ├── 每周 orphan files 清理 │ │
    │ │ └── 健康检查 + 告警 (Prometheus + Grafana) │ │
    │ └──────────────────────────────────────────────────────┘ │
    │ │
    └─────────────────────────────────────────────────────────────────────┘

    16.2 多环境管理

    多环境架构:

    ┌──────────────┐ ┌──────────────┐ ┌──────────────┐
    │ DEV 环境 │ │ STG 环境 │ │ PROD 环境 │
    │ │ │ │ │ │
    │ MinIO │ │ S3 │ │ S3 │
    │ Polaris │ │ Polaris │ │ Polaris HA │
    │ Spark │ │ Spark │ │ Spark HA │
    │ (单节点) │ │ (小规模) │ │ (生产集群) │
    │ │ │ │ │ │
    │ 自由实验 │ │ 验证流程 │ │ 严格管控 │
    └──────┬───────┘ └──────┬───────┘ └──────┬───────┘
    │ │ │
    └──────────────────┼──────────────────┘

    Git 版本管理:
    ├── catalog definitions
    ├── schema definitions
    ├── pipeline definitions
    └── maintenance schedules

    16.3 灾备方案

    灾备架构:

    主集群 (us-east-1) 灾备集群 (us-west-2)
    ┌──────────────────┐ ┌──────────────────┐
    │ S3 Primary │ ──复制──→│ S3 Replica │
    │ (S3 CRR) │ │ (Cross-Region) │
    │ │ │ │
    │ Polaris Primary │ ──同步──→│ Polaris Replica │
    │ PostgreSQL │ (pg流复制)│ PostgreSQL │
    │ │ │ │
    │ 活跃写入 │ │ 只读 (RTO < 1h) │
    └──────────────────┘ └──────────────────┘

    恢复流程:
    1. S3 跨区复制保证数据副本
    2. Polaris 元数据通过 PostgreSQL 流复制保持同步
    3. 灾备集群指向 Replica S3 + Replica Polaris
    4. Time Travel 提供额外数据恢复能力

    16.4 版本回滚

    — 生产环境版本回滚(利用 Time Travel)
    — Step 1: 查看快照历史,找到回滚目标
    SELECT
    snapshot_id,
    operation,
    committed_at,
    summary
    FROM demo.analytics.orders.snapshots
    ORDER BY committed_at DESC
    LIMIT 20;

    — Step 2: 回滚到指定快照
    — Spark SQL:
    CALL demo.system.rollback_to_snapshot(
    table => 'analytics.orders',
    snapshot_id => 1089405398375030
    );

    — PyIceberg:
    — table.manage_snapshots().rollback_to(1089405398375030).commit()

    — Step 3: 验证回滚结果
    SELECT COUNT(*) FROM demo.analytics.orders;

    16.5 成本优化

    优化方向具体措施预期节省
    存储成本 ZSTD 压缩 + S3 生命周期 40~60%
    计算成本 分布式文件规划 + 分区裁剪 50~80% 查询加速
    元数据成本 定期清理旧快照和孤立文件 避免元数据膨胀
    网络成本 同区域部署(计算和存储同 Region) 减少跨区域流量
    请求成本 增大目标文件大小(512MB) 减少 S3 GET/LIST 请求

    17. 实战项目

    17.1 项目一:电商订单数据湖

    目标:从 MySQL CDC 实时同步订单数据到 Iceberg,支持 OLAP 查询

    ┌──────────────────────────────────────────────────────────────┐
    │ 电商订单数据湖架构 │
    ├──────────────────────────────────────────────────────────────┤
    │ │
    │ MySQL (orders) │
    │ │ │
    │ Debezium CDC │
    │ │ │
    │ Kafka (mysql.orders.binlog) │
    │ │ │
    │ Flink Job (MERGE INTO) │
    │ │ │
    │ Iceberg Table (v3, MOR) │
    │ ├── bronze.orders_raw (原始 CDC 数据) │
    │ ├── silver.orders_clean (清洗后,去重) │
    │ └── gold.orders_summary (按天聚合报表) │
    │ │ │
    │ ┌───┴────────────────┐ │
    │ │ │ │
    │ Trino (BI查询) Spark (批量分析) │
    │ │
    └──────────────────────────────────────────────────────────────┘

    # —- Step 1: 创建分层表 —-
    from pyiceberg.catalog import load_catalog
    from pyiceberg.schema import Schema
    from pyiceberg.types import *

    catalog = load_catalog("demo")
    catalog.create_namespace_if_not_exists("ecommerce")

    # Bronze 层:原始 CDC 数据
    bronze_schema = Schema(
    NestedField(1, "order_id", LongType(), required=True),
    NestedField(2, "customer_id", LongType(), required=True),
    NestedField(3, "product_id", LongType(), required=True),
    NestedField(4, "quantity", IntegerType(), required=True),
    NestedField(5, "unit_price", DecimalType(10, 2)),
    NestedField(6, "total_amount", DecimalType(10, 2)),
    NestedField(7, "status", StringType()),
    NestedField(8, "created_at", TimestamptzType()),
    NestedField(9, "updated_at", TimestamptzType()),
    NestedField(10, "_cdc_op", StringType()), # INSERT/UPDATE/DELETE
    NestedField(11, "_cdc_ts", TimestamptzType()), # CDC 时间戳
    )

    bronze_table = catalog.create_table(
    identifier="ecommerce.bronze_orders",
    schema=bronze_schema,
    properties={
    "format-version": "3",
    "write.delete.mode": "merge-on-read",
    "write.update.mode": "merge-on-read",
    "write.parquet.compression-codec": "zstd",
    }
    )

    # Silver 层:清洗去重后的订单数据
    silver_schema = Schema(
    NestedField(1, "order_id", LongType(), required=True),
    NestedField(2, "customer_id", LongType(), required=True),
    NestedField(3, "product_id", LongType(), required=True),
    NestedField(4, "quantity", IntegerType(), required=True),
    NestedField(5, "unit_price", DecimalType(10, 2)),
    NestedField(6, "total_amount", DecimalType(10, 2)),
    NestedField(7, "status", StringType()),
    NestedField(8, "order_date", DateType()),
    NestedField(9, "created_at", TimestamptzType()),
    NestedField(10, "updated_at", TimestamptzType()),
    )

    silver_table = catalog.create_table(
    identifier="ecommerce.silver_orders",
    schema=silver_schema,
    partition_spec=..., # 按 order_date 天分区 + customer_id 分桶
    properties={
    "format-version": "3",
    "write.parquet.compression-codec": "zstd",
    "write.target-file-size-bytes": "536870912",
    }
    )

    — —- Step 2: Spark SQL 生成 Gold 层汇总 —-
    — 从 Silver 层聚合生成 Gold 层
    CREATE TABLE demo.ecommerce.gold_daily_summary USING iceberg
    TBLPROPERTIES ('format-version' = '3') AS
    SELECT
    order_date,
    COUNT(DISTINCT order_id) AS total_orders,
    COUNT(DISTINCT customer_id) AS unique_customers,
    SUM(total_amount) AS total_revenue,
    AVG(total_amount) AS avg_order_value,
    SUM(CASE WHEN status = 'completed' THEN 1 ELSE 0 END) AS completed_orders,
    SUM(CASE WHEN status = 'cancelled' THEN 1 ELSE 0 END) AS cancelled_orders
    FROM demo.ecommerce.silver_orders
    GROUP BY order_date;

    — —- Step 3: Trino 查询 —-
    — 每日订单趋势
    SELECT
    order_date,
    total_orders,
    total_revenue,
    total_revenue / total_orders AS avg_order_value
    FROM demo.ecommerce.gold_daily_summary
    WHERE order_date >= DATE '2024-01-01'
    ORDER BY order_date;

    — 客户消费排名
    SELECT
    customer_id,
    COUNT(*) AS order_count,
    SUM(total_amount) AS total_spent
    FROM demo.ecommerce.silver_orders
    WHERE order_date >= DATE '2024-01-01'
    GROUP BY customer_id
    ORDER BY total_spent DESC
    LIMIT 100;

    17.2 项目二:用户行为分析平台

    # 实时行为数据入湖 + 离线分析
    import pyarrow as pa
    from pyiceberg.catalog import load_catalog
    from datetime import datetime

    catalog = load_catalog("demo")

    # 创建行为数据表(带 VARIANT 列用于灵活的属性存储)
    schema = Schema(
    NestedField(1, "event_id", LongType(), required=True),
    NestedField(2, "user_id", LongType(), required=True),
    NestedField(3, "event_type", StringType(), required=True),
    NestedField(4, "page_url", StringType()),
    NestedField(5, "properties", VariantType()), # VARIANT 类型!
    NestedField(6, "event_time", TimestamptzType(), required=True),
    NestedField(7, "session_id", StringType()),
    NestedField(8, "device_type", StringType()),
    NestedField(9, "country", StringType()),
    )

    table = catalog.create_table(
    identifier="analytics.user_events",
    schema=schema,
    properties={
    "format-version": "3",
    "write.parquet.compression-codec": "zstd",
    }
    )

    # 模拟写入行为数据
    batch_data = pa.table({
    "event_id": [1, 2, 3],
    "user_id": [1001, 1002, 1001],
    "event_type": ["page_view", "click", "page_view"],
    "page_url": ["/home", "/products/123", "/about"],
    "properties": [
    '{"browser": "Chrome", "resolution": "1920×1080"}',
    '{"button": "add_to_cart", "price": 29.99}',
    '{"browser": "Safari", "resolution": "1440×900"}',
    ],
    "event_time": [datetime.now()] * 3,
    "session_id": ["sess_abc", "sess_def", "sess_abc"],
    "device_type": ["desktop", "mobile", "desktop"],
    "country": ["US", "UK", "US"],
    })

    table.append(batch_data)

    # 离线分析:用户行为漏斗
    spark.sql("""
    WITH event_funnel AS (
    SELECT
    user_id,
    session_id,
    event_type,
    event_time,
    ROW_NUMBER() OVER (PARTITION BY user_id, session_id ORDER BY event_time) AS step
    FROM demo.analytics.user_events
    WHERE event_time >= CURRENT_DATE – INTERVAL 7 DAYS
    )
    SELECT
    event_type,
    COUNT(DISTINCT user_id) AS unique_users,
    COUNT(DISTINCT session_id) AS total_sessions
    FROM event_funnel
    GROUP BY event_type
    ORDER BY total_sessions DESC
    """
    ).show()

    17.3 项目三:跨引擎联邦查询

    ┌──────────────────────────────────────────────────────────────┐
    │ 跨引擎联邦查询 │
    ├──────────────────────────────────────────────────────────────┤
    │ │
    │ ┌──────────┐ 写入 ┌──────────────────────┐ │
    │ │ Spark │ ──────────→ │ Iceberg Table │ │
    │ │ (ETL) │ │ (共享存储层) │ │
    │ └──────────┘ └──────┬───────┬───────┘ │
    │ │ │ │
    │ 读取 │ │ 读取 │
    │ ┌────────────┘ └────────┐ │
    │ ↓ ↓ │
    │ ┌──────────────┐ ┌──────────────┐ │
    │ │ Trino │ │ Snowflake │ │
    │ │ (OLAP 查询) │ │ (数据共享) │ │
    │ └──────────────┘ └──────────────┘ │
    │ │
    │ 同一份数据,不同引擎各自优化 │
    │ │
    └──────────────────────────────────────────────────────────────┘

    # Spark 写入(ETL 处理)
    spark.sql("""
    INSERT OVERWRITE demo.analytics.customer_360
    SELECT
    c.customer_id,
    c.name,
    c.email,
    COUNT(DISTINCT o.order_id) AS total_orders,
    SUM(o.total_amount) AS lifetime_value,
    MAX(o.order_date) AS last_order_date,
    DATEDIFF(CURRENT_DATE, MAX(o.order_date)) AS days_since_last_order
    FROM demo.ecommerce.silver_customers c
    LEFT JOIN demo.ecommerce.silver_orders o ON c.customer_id = o.customer_id
    GROUP BY c.customer_id, c.name, c.email
    """
    )

    — Trino 查询(BI 报表)
    SELECT
    CASE
    WHEN days_since_last_order <= 30 THEN 'Active'
    WHEN days_since_last_order <= 90 THEN 'At Risk'
    ELSE 'Churned'
    END AS customer_segment,
    COUNT(*) AS customer_count,
    AVG(lifetime_value) AS avg_ltv
    FROM demo.analytics.customer_360
    GROUP BY 1
    ORDER BY avg_ltv DESC;

    — Snowflake 查询(数据共享给外部合作伙伴)
    SELECT
    customer_segment,
    COUNT(*) AS customers,
    SUM(lifetime_value) AS total_value
    FROM demo.analytics.customer_360
    GROUP BY customer_segment;

    18. 2025-2026 技术趋势

    18.1 Iceberg v4 规划

    v4 目前处于活跃设计阶段,尚未发布,预计在 2027 年或之后才会落地:

    特性状态说明
    Single-file Commits 积极开发中 解决频繁提交的元数据开销
    Adaptive Root Manifest 积极开发中 动态元数据树,大规模表性能优化
    Relative Paths 已基本确定 表可移植,支持灾难恢复和表迁移
    Content Statistics 已基本确定 更丰富的内容统计信息
    Partition Tuple 改进 讨论中 分区元数据结构的优化
    列级更新优化 讨论中 依赖 Parquet 的 logical-file 概念

    Single-file Commits 要解决的问题:

    当前 (v3):
    每次提交 → 新 metadata.json + 新 manifest list + 新 manifest
    → 元数据文件数量线性增长
    → 频繁提交时(如 Flink 30s checkpoint)元数据管理成本高

    v4 目标:
    多次提交合并到单个元数据文件中
    → 减少元数据文件数量
    → 降低频繁提交的开销
    → 特别适合流式写入场景

    18.2 多语言客户端生态

    Iceberg 多语言客户端发展状况(2026 年中):

    ┌────────────────────────────────────────────┐
    │ Java (参考实现) │
    │ 版本: 1.11.0 │
    │ 状态: 最成熟,功能最完整 │
    │ 测试: ~10,000 Spark 集成测试 │
    └────────────────────────────────────────────┘

    ┌────────────────────────────────────────────┐
    │ Python (PyIceberg) │
    │ 版本: 0.11.0 │
    │ 状态: 成熟,380+ PR,50+ 贡献者 │
    │ 特点: 纯 Python,不依赖 JVM │
    └────────────────────────────────────────────┘

    ┌────────────────────────────────────────────┐
    │ Rust (iceberg-rust) │
    │ 版本: 0.10.0 │
    │ 状态: 快速成熟,254 PR,40 贡献者 │
    │ 特点: 原生高性能,提供 Python 绑定(pyiceberg-core)│
    └────────────────────────────────────────────┘

    ┌────────────────────────────────────────────┐
    │ Go (iceberg-go) │
    │ 版本: 0.6.0 │
    │ 状态: 积极开发,200 PR,40 贡献者 │
    │ 特点: Go 生态集成 │
    └────────────────────────────────────────────┘

    ┌────────────────────────────────────────────┐
    │ C++ (iceberg-cpp) │
    │ 版本: 0.3.0 │
    │ 状态: 早期阶段,140+ PR,23 贡献者 │
    │ 特点: 嵌入式场景,Arrow 生态 │
    └────────────────────────────────────────────┘

    Rust 客户端与 DataFusion Comet 的协同:

    Rust 客户端通过与 Apache DataFusion Comet 的集成,正在加速 Spark 查询性能。iceberg-rust 利用 Iceberg Java 的近万个 Spark 测试用例作为差分测试框架,甚至帮助发现了 Java 实现中的 Bug。

    18.3 Iceberg + AI

    Iceberg 在 AI/ML 场景的应用:

    ┌─────────────────────────────────────────────────────┐
    │ AI 工作负载与 Iceberg │
    ├─────────────────────────────────────────────────────┤
    │ │
    │ 1. 特征存储 (Feature Store) │
    │ ├── 离线特征:Iceberg 表存储历史特征 │
    │ ├── 在线特征:Flink 实时更新 Iceberg 表 │
    │ └── 特征版本:Time Travel 管理特征版本 │
    │ │
    │ 2. 向量存储 (Vector Storage) │
    │ ├── Iceberg 表存储 embedding 向量 │
    │ ├── 通过 VARIANT 类型存储元数据 │
    │ └── 配合外部向量索引(如 Weaviate) │
    │ │
    │ 3. 训练数据管理 │
    │ ├── 训练数据集版本化(快照/标签/分支) │
    │ ├── 数据血缘追踪(Row Lineage) │
    │ └── 可重现性(Time Travel 回溯训练数据) │
    │ │
    │ 4. MLOps 集成 │
    │ ├── PyIceberg 与 MLflow/Feast 集成 │
    │ ├── 模型输入数据的 ACID 保证 │
    │ └── 模型推理结果写回 Iceberg │
    │ │
    └─────────────────────────────────────────────────────┘

    18.4 生态融合趋势

    2026 年湖格式生态格局:

    ┌─────────────────────────────────────────────────────┐
    │ │
    │ Iceberg ←──── 事实标准 ────→ 最广引擎支持 │
    │ ↑ ↑ │
    │ │ │ │
    │ Delta UniForm ──→ 兼容读取 Iceberg │
    │ ↑ │
    │ Databricks 生态 │
    │ │
    │ Hudi ←──── CDC 特化 ────→ XTable 导出 Iceberg │
    │ │
    │ Paimon ←── Flink 特化 ────→ Iceberg-compat 模式 │
    │ │
    │ 趋势:所有格式都在构建通往 Iceberg 的桥梁 │
    │ │
    └─────────────────────────────────────────────────────┘

    趋势说明影响
    Delta → Iceberg Delta UniForm 允许 Delta 表被 Iceberg 引擎原生读取 格式边界模糊化
    Hudi → Iceberg XTable 项目支持 Hudi 表导出为 Iceberg 格式 Hudi 用户可渐进迁移
    Paimon → Iceberg Paimon 的 Iceberg-compat 模式支持以 Iceberg 格式暴露 Flink 用户也能享受 Iceberg 生态
    Catalog 标准化 REST Catalog 成为统一标准,Apache Polaris 毕业 跨格式跨引擎治理
    Parquet 协同 Parquet Variant 与 Iceberg Variant 同步演进 文件格式与表格式深度协同

    18.5 未来展望

    2026-2028 路线图展望:

    2026 H2:
    ├── Iceberg 1.12.0 发布
    ├── 更多引擎完善 v3 支持
    ├── Polaris 成为 REST Catalog 标准实现
    └── Rust/Go 客户端进入生产可用

    2027:
    ├── Iceberg v4 规范草案
    ├── Single-file commits 落地
    ├── 多参数分区转换全面铺开
    └── AI 特征存储标准化

    2028+:
    ├── v4 正式发布
    ├── 格式融合进一步深入
    ├── 湖仓一体成为企业标配
    └── 实时分析延迟降至亚秒级

    🏆 面试加分提示:2026 年选择 Iceberg 已经不需要论证"为什么要用 Iceberg",问题变成了"怎么用好 Iceberg"。Iceberg Summit 2026 的 70+ 场技术分享中没有一场在讨论"是否采用",全部聚焦于"如何改进"。这种跨厂商共识是 Iceberg 最大的护城河——来自 Google、Apple、Snowflake、Databricks、Microsoft、Netflix、LinkedIn 的工程师在同一张桌子上讨论设计决策。

    ⚠️ 生产踩坑提醒(全局):

  • 永远不要跳过表维护——Compaction、Snapshot 清理、Orphan 清理是生产必备
  • 选择 REST Catalog——它是 2026 年的事实标准,提供凭证代理和 RBAC
  • 从小规模开始——先在 DEV 环境验证,再推到 STG/PROD
  • 监控先行——部署 Iceberg 前,先搭建好监控和告警
  • 拥抱 v3——Deletion Vectors 和 Row Lineage 对 CDC 场景是变革性的
  • 保持数据文件在开放格式——Parquet 是你最深的保险策略

  • 补充专题


    补充一:Iceberg File Format API(2026 年新里程碑)

    2026 年 2 月,Apache Iceberg 社区宣布 File Format API 正式定稿,这是一个重大的架构里程碑。它使文件格式变得可插拔、一致且引擎无关。

    为什么需要 File Format API?

    在 Iceberg 1.11.0 之前,文件格式(Parquet、ORC、Avro)的读写逻辑散布在 Iceberg Java 代码的各个角落。添加新的文件格式需要修改大量内部代码,且不同引擎对文件格式的处理方式不一致。File Format API 定义了一个统一的接口层:

    File Format API 架构:

    ┌─────────────────────────────────────────────┐
    │ Iceberg Table API │
    │ (Catalog / Table / Scan / Write) │
    ├─────────────────────────────────────────────┤
    │ File Format API (新!) │
    │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
    │ │ Parquet │ │ ORC │ │ Avro │ │
    │ │ Reader │ │ Reader │ │ Reader │ │
    │ │ Writer │ │ Writer │ │ Writer │ │
    │ └─────────┘ └─────────┘ └─────────┘ │
    │ ┌─────────┐ ┌─────────┐ │
    │ │ Custom │ │ Future │ │
    │ │ Format │ │ Formats │ │
    │ └─────────┘ └─────────┘ │
    ├─────────────────────────────────────────────┤
    │ 存储层 (S3 / HDFS / ADLS) │
    └─────────────────────────────────────────────┘

    File Format API 的关键价值:

    价值说明
    可插拔 第三方可以实现自定义文件格式并接入 Iceberg
    一致性 所有引擎通过相同 API 读写文件格式,消除引擎差异
    引擎无关 文件格式的实现独立于计算引擎
    测试复用 Iceberg Java 的近万个 Spark 测试可以验证任何文件格式实现
    Rust 协同 iceberg-rust 可以通过相同 API 提供文件格式支持

    补充二:DuckDB + Iceberg:桌面级湖仓分析

    DuckDB 是 2026 年数据工程师最热门的本地分析工具之一。通过 Iceberg REST Catalog 扩展,DuckDB 可以直接查询 Iceberg 表,无需 Spark 或 Flink 集群。

    DuckDB 连接 Iceberg 配置

    — 安装并加载 Iceberg 扩展
    INSTALL iceberg;
    LOAD iceberg;

    — 配置 REST Catalog 连接
    CREATE SECRET iceberg_secret (
    TYPE ICEBERG,
    CATALOG_URL 'http://localhost:8181',
    WAREHOUSE 's3://warehouse',
    ENDPOINT 'http://localhost:9000',
    ACCESS_KEY_ID 'admin',
    SECRET_ACCESS_KEY 'localpassword',
    USE_SSL false
    );

    — 查询 Iceberg 表
    SELECT
    customer_id,
    COUNT(*) AS order_count,
    SUM(total_amount) AS total_spent
    FROM iceberg_catalog.analytics.orders
    WHERE order_date >= '2024-01-01'
    GROUP BY customer_id
    ORDER BY total_spent DESC
    LIMIT 10;

    — Time Travel 查询
    SELECT * FROM iceberg_catalog.analytics.orders
    AT SNAPSHOT 1089405398375037;

    DuckDB + PyIceberg + Polars 组合

    2026 年最流行的本地数据分析组合是:PyIceberg(Catalog 操作) + Polars(DataFrame 计算) + DuckDB(SQL 查询)。

    from pyiceberg.catalog import load_catalog
    import polars as pl

    # 使用 PyIceberg 获取数据
    catalog = load_catalog("local", type="rest", uri="http://localhost:8181")
    table = catalog.load_table("analytics.orders")

    # 用 PyIceberg 读取为 Arrow 表
    arrow_table = table.scan(
    row_filter="order_date >= '2024-01-01'",
    selected_fields=("customer_id", "total_amount", "order_date")
    ).to_arrow()

    # 转换为 Polars DataFrame 进行分析
    df = pl.from_arrow(arrow_table)

    result = (
    df.group_by("customer_id")
    .agg([
    pl.count().alias("order_count"),
    pl.sum("total_amount").alias("total_spent"),
    pl.max("order_date").alias("last_order"),
    ])
    .sort("total_spent", descending=True)
    .head(100)
    )

    print(result)

    这个组合的优势在于:

    • 零基础设施:不需要 Spark 集群、Flink 集群或任何分布式系统
    • 内存计算:DuckDB 和 Polars 都是内存列式引擎,速度极快
    • 完整 Catalog 支持:通过 REST Catalog 可以访问生产环境的数据
    • 开发体验:比启动 Spark Session 快 100 倍

    补充三:流式入湖方案对比深度分析

    对于需要将实时数据写入 Iceberg 的团队,选择正确的入湖方案至关重要。以下是 2026 年主流方案的深度对比:

    三种流式入湖方案详细对比

    维度Apache FlinkSpark Structured StreamingKafka Connect Iceberg Sink
    延迟 秒级(30s checkpoint) 分钟级(1~5min trigger) 分钟级(5min commit)
    写入模式 每个 checkpoint 一次提交 每个 trigger 一次提交 每个 commit interval 一次提交
    文件产生速率 高(32 files/min @16p) 中(适度 batching) 低(gentler on table)
    Exactly-Once ✅ 通过 checkpoint 协调 ✅ 通过微批次 ⚠️ at-least-once
    Upsert/MERGE ✅ 完整支持 ✅ 通过 foreachBatch ❌ 仅追加
    Schema Evolution ⚠️ 需要停流处理 ✅ 在线演进 ✅ 配合 Schema Registry
    CDC 支持 ✅ 原生 Debezium ✅ 通过 foreachBatch ⚠️ 有限
    运维复杂度 高(checkpoint 管理) 中(微批次简单) 低(配置式)
    最佳场景 低延迟 CDC 入湖 已有 Spark 集群 简单追加管道
    生产建议 最推荐 团队已用 Spark 时 最简单但功能有限

    Flink 入湖的文件管理挑战

    Flink 30s checkpoint 的文件管理数学:

    假设条件:
    ├── 16 个写入并行度
    ├── 30 秒 checkpoint 间隔
    └── 每个 checkpoint 每个并行度产生 2 个文件

    每分钟文件产生:
    16 parallelism × 2 files × 2 checkpoints = 64 files/min

    每小时:64 × 60 = 3,840 files/hour
    每天:3,840 × 24 = 92,160 files/day
    每周:92,160 × 7 = 645,120 files/week!

    这就是为什么 compaction 不是"可选的",而是"必须的"。
    没有自动 compaction 的流式 Iceberg 表,一周内就会变得不可用。

    推荐方案:Flink 入湖 + LakeOps 自动维护

    推荐架构:

    ┌─────────┐ ┌────────┐ ┌──────────┐ ┌──────────┐
    │ Kafka │───→│ Flink │───→│ Iceberg │←───│ LakeOps │
    │ (源) │ │ (写入) │ │ Table │ │ (自动 │
    └─────────┘ └────────┘ │ (v3 MoR) │ │ 维护) │
    └──────────┘ └──────────┘

    ┌───────┼───────┐
    │ │ │
    Compaction Expire Orphan
    (持续) (每日) (每周)

    LakeOps 是一个自治的 Iceberg 表管理平台:
    ├── 感知每个表的写入模式
    ├── 自动触发 compaction(基于文件数量阈值)
    ├── 自动管理快照过期
    ├── 自动清理孤立文件
    └── 提供每个表的健康仪表板


    补充四:Iceberg 多语言客户端开发实战

    2026 年,Iceberg 不再只是 Java 的天下。Rust、Go、C++ 客户端已经可以在生产环境中使用。

    Rust 客户端(iceberg-rust)快速开始

    # Cargo.toml
    [dependencies]
    iceberg = "0.10"
    iceberg-catalog-rest = "0.10"
    iceberg-datafusion = "0.10"

    // Rust: 连接 REST Catalog 并读取 Iceberg 表
    use iceberg::catalog::rest::RestCatalog;
    use iceberg::catalog::Catalog;

    #[tokio::main]
    async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 配置 REST Catalog
    let catalog = RestCatalog::builder()
    .uri("http://localhost:8181")
    .warehouse("s3://warehouse")
    .build()?;

    // 加载表
    let table = catalog
    .load_table(&["analytics", "orders"].into())
    .await?;

    // 扫描数据
    let scan = table.scan().build()?;
    let batches = scan.to_datafusion_batches().await?;

    for batch in batches {
    println!("Rows: {}", batch.num_rows());
    }

    Ok(())
    }

    Go 客户端(iceberg-go)快速开始

    // Go: 连接 Iceberg 表
    package main

    import (
    "context"
    "fmt"

    "github.com/apache/iceberg-go"
    "github.com/apache/iceberg-go/catalog/rest"
    )

    func main() {
    ctx := context.Background()

    // 创建 REST Catalog 客户端
    cat, err := rest.NewCatalog(
    ctx,
    "demo",
    "http://localhost:8181",
    rest.WithWarehouse("s3://warehouse"),
    )
    if err != nil {
    panic(err)
    }

    // 列出命名空间
    namespaces, err := cat.ListNamespaces(ctx)
    if err != nil {
    panic(err)
    }

    for _, ns := range namespaces {
    fmt.Printf("Namespace: %v\\n", ns)
    }
    }

    多语言客户端成熟度对比(2026年中)

    能力Java (1.11.0)Python (0.11.0)Rust (0.10.0)Go (0.6.0)C++ (0.3.0)
    REST Catalog ⚠️
    Hive Catalog
    Parquet 读写 ⚠️ ⚠️
    ORC 读写 ⚠️
    Schema Evolution ⚠️ ⚠️
    v3 Deletion Vectors ⚠️
    v3 Row Lineage
    v3 VARIANT ⚠️
    Table Scan ⚠️ ⚠️
    文件 Compaction
    生产就绪度 🟢 🟢 🟡 🟡 🔴

    补充五:数据湖设计模式与最佳实践

    模式一:Medallion Architecture(奖牌架构)

    Medallion Architecture 数据分层:

    ┌───────────────────────────────────────────────────────────┐
    │ │
    │ Bronze Layer (原始层) │
    │ ├── 原始数据的精确副本 │
    │ ├── 不做任何转换 │
    │ ├── Schema-on-Read(灵活解析) │
    │ └── 用于数据溯源和审计 │
    │ │
    │ Silver Layer (清洗层) │
    │ ├── 数据清洗、去重、类型转换 │
    │ ├── Schema 标准化 │
    │ ├── 业务规则验证 │
    │ └── 用于下游分析和 ML │
    │ │
    │ Gold Layer (聚合层) │
    │ ├── 业务聚合和汇总 │
    │ ├── 面向特定业务场景 │
    │ ├── 高度优化(排序、分区、compaction) │
    │ └── 用于 BI 报表和实时查询 │
    │ │
    └───────────────────────────────────────────────────────────┘

    Iceberg 实现建议:
    ├── 每个层级使用独立的命名空间:bronze / silver / gold
    ├── Bronze 用追加模式(append-only),不做删除
    ├── Silver 用 MOR 模式(支持 CDC 更新)
    └── Gold 用 COW 模式(查询性能优先)

    模式二:Change Data Capture(变更数据捕获)

    CDC 入湖的三种模式对比:

    模式 A: 追加原始 CDC 日志
    ┌─────────┐ ┌──────────────────────────────┐
    │ Debezium│───→│ Bronze: raw_cdc_events │
    │ │ │ (order_id, op, before, after) │
    └─────────┘ │ 追加模式,保留完整历史 │
    └──────────────────────────────┘

    模式 B: 实时 Upsert 到目标表
    ┌─────────┐ ┌──────────────────────────────┐
    │ Debezium│───→│ Silver: orders │
    │ + Flink │ │ (order_id, amount, status) │
    │ (MERGE) │ │ 主键表,实时更新 │
    └─────────┘ └──────────────────────────────┘

    模式 C: 混合模式(推荐)
    ┌─────────┐ ┌───────────┐ ┌───────────┐
    │ Debezium│───→│ Bronze │───→│ Silver │
    │ │ │ (raw CDC) │ │ (Upsert) │
    └─────────┘ └───────────┘ └───────────┘
    保留原始记录 提供干净的
    最新状态

    模式三:Time Travel 数据回滚

    生产环境数据回滚决策树:

    发现数据错误

    ├── 错误在 24 小时内?
    │ ├── 是 → Time Travel 查询旧快照
    │ │ └── INSERT OVERWRITE 覆盖受影响分区
    │ └── 否 → 快照可能已过期
    │ └── 从备份恢复或重新处理

    ├── 错误范围?
    │ ├── 单个分区 → 只覆盖该分区(快速)
    │ ├── 多个分区 → 批量覆盖受影响分区
    │ └── 全表 → INSERT OVERWRITE 全表(谨慎)

    └── 回滚后验证
    ├── 数据量检查
    ├── 聚合值检查
    └── 与源系统对账

    模式四:数据共享与联邦查询

    跨组织数据共享架构:

    组织 A (数据提供方) 组织 B (数据消费方)
    ┌──────────────────┐ ┌──────────────────┐
    │ Iceberg Tables │ │ Iceberg Tables │
    │ (Gold Layer) │ │ │
    └────────┬─────────┘ └────────┬─────────┘
    │ │
    │ REST Catalog API │
    │ (Apache Polaris) │
    └────────────┬───────────────────┘

    ┌─────────┴─────────┐
    │ 数据共享层 │
    │ ├── 凭证代理 │
    │ ├── 细粒度权限 │
    │ └── 审计日志 │
    └───────────────────┘

    实现方式(2026年):
    ├── Snowflake 的 Iceberg Tables 数据共享
    ├── Apache Polaris 的跨组织 REST Catalog
    ├── AWS Data Exchange + Glue + Iceberg
    └── 通用方案:REST Catalog + 凭证代理 + RBAC


    补充六:Iceberg 性能基准参考

    以下是社区和生产环境中常见的 Iceberg 性能参考数据(基于公开基准测试和报告):

    操作数据量引擎耗时备注
    全表扫描 1TB (10亿行) Spark 4.0 + Iceberg ~2-3 min ZSTD 压缩,128 核
    分区过滤查询 1TB → 10GB Spark 4.0 + Iceberg ~15-30 sec 命中分区裁剪
    Point Lookup 单行 Trino + Iceberg ~100-500 ms 无排序优化
    排序表查询 1TB → 范围 Trino + Iceberg ~5-10 sec 排序优化后
    INSERT 追加 1000万行 Flink ~30 sec 单个 checkpoint
    MERGE INTO 100万行匹配 Spark ~2-5 min v3 MoR 模式
    Compaction 10000 小文件→500 Spark ~10-30 min BinPack 策略
    快照过期 1000 快照清理 Spark ~1-5 min stream_results=true

    影响性能的关键因素排名

    性能影响因子(从大到小):

    1. 数据文件大小和数量 ★★★★★ → 影响最大
    ├── 太小:文件打开开销巨大
    └── 太大:并行度不足

    2. 分区设计 ★★★★☆ → 分区裁剪效果
    ├── 好的分区:跳过 99% 数据
    └── 差的分区:全表扫描

    3. 排序顺序 ★★★★☆ → 范围扫描加速
    ├── 按查询键排序:高效范围扫描
    └── 无排序:随机 IO

    4. 压缩算法 ★★★☆☆ → IO 开销
    ├── ZSTD:最佳平衡
    └── Snappy:最快解压

    5. Catalog 延迟 ★★☆☆☆ → 元数据获取
    ├── REST Catalog:HTTP 延迟
    └── Hive Catalog:Thrift 延迟

    6. 网络带宽 ★★☆☆☆ → 数据传输
    ├── 同 Region:低延迟
    └── 跨区域:高延迟


    附录:深度扩展内容


    前置知识:从 Hive 到 Iceberg 的演进之路

    要真正理解 Iceberg 的价值,需要先理解它的前辈——Apache Hive 的局限性。Hive 是大数据时代的开创性项目,它将 MapReduce 引入数据仓库领域,使 SQL 查询可以运行在 HDFS 上的海量数据上。但随着数据规模的增长和业务需求的变化,Hive 的架构缺陷日益暴露:

    Hive 的六大痛点

  • 无事务保证:两个 Spark 作业同时写入同一张 Hive 表,数据可能损坏。没有锁机制,没有隔离级别,写入结果取决于文件系统的最终一致性。这在 Netflix 这样的生产环境中是不可接受的。

  • 分区方案固化:Hive 的分区直接映射到 HDFS 目录结构。如果最初按 date 分区,后来需要改为按 date_hour 分区,就必须重写 PB 级的数据。这个过程可能需要数天,且期间需要停止所有下游作业。

  • Schema 变更代价高昂:Hive 的 Schema 存储在 Metastore 中,按列位置映射。新增一列只在末尾追加,不能插入中间位置。如果业务需要调整列顺序,必须修改所有数据文件。

  • 小文件灾难:Spark 的每个 Task 产生一个输出文件。一个 1000 个 Map 任务的作业可能产生 1000 个文件。日积月累,一个分区可能有数百万个小文件。Hive 查询需要打开所有文件,性能急剧下降。

  • 无时间旅行:如果某个 ETL 作业写入了错误数据,无法回滚到之前的状态。只能手动定位错误数据、重新运行作业、覆盖文件。

  • 文件格式与分区绑定:文件格式(Parquet/ORC)的统计信息不被查询引擎充分利用。Hive 查询需要扫描所有文件,即使某些文件根本不包含匹配数据。

  • Iceberg 如何解决这些问题

    Hive 痛点Iceberg 解决方案核心机制
    无事务 ACID 事务 原子性元数据交换 + MVCC
    分区固化 隐藏分区 + 分区演进 分区转换函数存储在元数据中
    Schema 变更 Schema Evolution 基于列 ID 的 Schema 追踪
    小文件 文件规划 + Compaction Manifest 文件过滤 + 定期合并
    无时间旅行 Snapshot 快照 每次写入产生新快照
    统计信息利用 列统计下推 Manifest 中的 min/max/null_count

    理解了这个背景,你就能明白为什么 Iceberg 的每一个设计决策都是必要的。它不是凭空发明的——它是从 Netflix 数十 PB 数据的实际痛点中生长出来的。


    深度理解:Iceberg 的元数据层级与查询路径

    当一个查询到达 Iceberg 表时,引擎经历以下路径来获取数据:

    查询执行路径(以 Trino 为例):

    用户提交查询:
    SELECT * FROM orders WHERE customer_id = 1001 AND order_date = '2024-01-15'

    Step 1: Catalog 解析
    ┌──────────────────────────────────────────┐
    │ Trino → REST Catalog: "加载 orders 表" │
    │ Catalog 返回: │
    │ ├── 当前 metadata.json 的位置 │
    │ ├── 当前 schema │
    │ └── 当前 partition spec │
    │ 耗时: ~10-50ms │
    └──────────────────────────────────────────┘

    Step 2: 读取 Metadata JSON
    ┌──────────────────────────────────────────┐
    │ 从 S3 读取 metadata.json │
    │ 获取当前快照的 manifest list 位置 │
    │ 耗时: ~10-100ms (S3 GET) │
    └──────────────────────────────────────────┘

    Step 3: 读取 Manifest List
    ┌──────────────────────────────────────────┐
    │ 从 S3 读取 snap-*.avro │
    │ 获取该快照的所有 manifest file 列表 │
    │ 每个 manifest 条目包含分区统计信息 │
    │ → 过滤掉不匹配 order_date='2024-01-15' │
    │ 的 manifest │
    │ 耗时: ~10-50ms │
    └──────────────────────────────────────────┘

    Step 4: 读取 Manifest Files
    ┌──────────────────────────────────────────┐
    │ 从 S3 读取匹配的 manifest files │
    │ 每个 data file 条目包含: │
    │ ├── partition values → 进一步分区裁剪 │
    │ ├── lower_bounds → 文件级 min 过滤 │
    │ ├── upper_bounds → 文件级 max 过滤 │
    │ └── null_value_counts → 空值统计 │
    │ → 过滤掉 customer_id 范围不匹配的文件 │
    │ 耗时: ~50-200ms (多个 manifest) │
    └──────────────────────────────────────────┘

    Step 5: 规划文件切片
    ┌──────────────────────────────────────────┐
    │ 将匹配的数据文件分组为 splits │
    │ 每个 split 约 128MB │
    │ 分配给不同的查询任务 │
    └──────────────────────────────────────────┘

    Step 6: 读取数据文件
    ┌──────────────────────────────────────────┐
    │ 从 S3 读取 Parquet 数据文件 │
    │ 每个文件: │
    │ ├── 读取 Footer → Row Group 统计 │
    │ ├── 过滤不匹配的 Row Group │
    │ ├── 读取匹配的 Row Group │
    │ └── 列式解码 → 过滤条件 │
    │ 耗时: 取决于匹配数据量 │
    └──────────────────────────────────────────┘

    这个六步路径展示了 Iceberg 的核心优势:用元数据代替数据扫描。在 Step 3 和 Step 4 中,引擎可以跳过 90%~99% 的数据文件,根本不需要打开它们。这就是为什么 Iceberg 的查询性能比直接扫描 S3 文件快数十倍。

    元数据缓存的重要性

    由于每次查询都需要读取 metadata.json → manifest list → manifest files,这些元数据文件的访问延迟直接影响查询性能。这就是为什么 Iceberg 引擎都实现了元数据缓存:

    缓存层缓存内容默认 TTL建议
    Metadata JSON 缓存 metadata.json 内容 1 分钟 表属性变更频繁时缩短
    Manifest List 缓存 snap-*.avro 内容 1 分钟 通常不需要调整
    Manifest File 缓存 manifest 条目 1 分钟 大表可增加
    文件统计缓存 lower/upper bounds 查询级 始终启用

    进阶概念:Iceberg 的写入路径详解

    理解写入路径有助于优化写入性能:

    写入路径(以 INSERT 为例):

    1. 引擎准备数据
    ├── 从上游读取数据(DataFrame / SQL 结果)
    ├── 按分区规范计算每条记录的分区值
    └── 按排序顺序排序数据(如果配置了 sort order)

    2. 写入数据文件
    ├── 每个写入任务产生一个或多个 Parquet 文件
    ├── 文件大小由 write.target-file-size-bytes 控制
    ├── 压缩由 write.parquet.compression-codec 控制
    └── 每个文件的统计信息(min/max/null_count)被记录

    3. 写入 Manifest File
    ├── 新产生的数据文件被记录到新的 manifest file
    ├── manifest file 记录每个数据文件的元信息
    └── manifest file 是 Avro 格式

    4. 写入 Manifest List
    ├── 新的 manifest list 引用所有 manifest files
    │ (包括新写入的和保持不变的历史 manifest)
    └── manifest list 也是 Avro 格式

    5. 写入 Metadata JSON
    ├── 新的 metadata.json 包含完整的表状态
    ├── 引用新的 manifest list
    ├── 记录新的快照信息
    └── 更新 current-snapshot-id

    6. 原子提交
    ├── 通过 Catalog 原子地将 metadata 指针切换到新版本
    ├── 成功 → 新快照可见
    └── 失败(冲突)→ 重试

    写入优化技巧总结

    技巧效果适用场景
    增大目标文件大小 减少文件数量,降低元数据开销 大表写入
    启用 ZSTD 压缩 减少存储空间和 IO 开销 通用推荐
    配置 hash 分布模式 均匀分布数据到分区 多分区写入
    增大写入并行度 加速大表写入 大表初始加载
    减少提交频率 减少元数据文件数量 流式写入
    使用 INSERT OVERWRITE 避免合并开销 分区级全量替换

    A.1 Iceberg 事务模型深度解析

    Iceberg 的事务模型是其区别于传统数据湖的核心特性之一。理解事务模型对于正确使用 Iceberg 至关重要。

    乐观并发控制(Optimistic Concurrency Control)

    Iceberg 事务提交流程:

    Writer 1 Writer 2
    ┌──────────┐ ┌──────────┐
    │ 读取当前 │ │ 读取当前 │
    │ snapshot │ │ snapshot │
    │ (N) │ │ (N) │
    └────┬─────┘ └────┬─────┘
    │ │
    │ 处理数据 │ 处理数据
    │ 写入新的 │ 写入新的
    │ data files │ data files
    │ │
    ┌────┴─────┐ ┌────┴─────┐
    │ 准备提交 │ │ 准备提交 │
    │ 新快照 │ │ 新快照 │
    │ (N+1) │ │ (N+1) │
    └────┬─────┘ └────┬─────┘
    │ │
    │ 尝试原子提交 │ 尝试原子提交
    │ CAS: N → N+1 │ CAS: N → N+1
    │ │
    ┌────┴─────┐ ┌────┴─────┐
    │ ✅ 成功 │ │ ❌ 冲突 │
    │ 当前快照 │ │ 重试: │
    │ = N+1 │ │ 读取 N+1 │
    └──────────┘ │ 重新规划 │
    │ 提交 N+2 │
    └──────────┘

    关键特性:
    ├── 读取永远不被阻塞(MVCC)
    ├── 写入通过 CAS(Compare-And-Swap)保证原子性
    ├── 冲突的 writer 自动重试
    └── 没有锁、没有阻塞、没有死锁

    多表事务(Multi-Table Transactions)

    REST Catalog(特别是 Apache Polaris)支持多表事务:

    # 多表事务示例(通过 Polaris REST Catalog)
    # 确保多张表要么同时成功,要么同时失败

    # 场景:订单表和库存表需要原子更新
    # 1. 创建新订单 → 订单表增加记录
    # 2. 扣减库存 → 库存表更新记录
    # 两个操作必须在同一个事务中完成

    # 使用 Spark 的多表事务 API
    spark.sql("""
    BEGIN TRANSACTION;

    INSERT INTO demo.ecommerce.orders
    VALUES (9999, 1001, '2024-01-20', 299.99, 'pending', now(), now());

    UPDATE demo.ecommerce.inventory
    SET quantity = quantity – 1
    WHERE product_id = 5001;

    COMMIT;
    """)

    冲突解决策略

    冲突类型说明解决方式
    写-写冲突 两个 writer 同时修改同一分区 后提交的 writer 重试
    写-读冲突 writer 写入时 reader 正在读 不影响(MVCC)
    Schema 冲突 两个 writer 同时修改 Schema 后提交的失败
    分区冲突 两个 writer 修改同一分区的分区规范 后提交的失败

    A.2 Iceberg 与 Parquet 深度协同

    Parquet 是 Iceberg 的默认和推荐数据文件格式。理解 Parquet 的内部结构有助于优化 Iceberg 查询性能。

    Parquet 文件内部结构

    Parquet 文件结构:

    ┌──────────────────────────────────────────────┐
    │ 4 bytes: Magic Number "PAR1" │
    ├──────────────────────────────────────────────┤
    │ Row Group 1 │
    │ ├── Column Chunk: order_id │
    │ │ ├── Page 1 (min=1, max=1000) │
    │ │ ├── Page 2 (min=1001, max=2000) │
    │ │ └── … │
    │ ├── Column Chunk: customer_id │
    │ │ ├── Page 1 │
    │ │ └── … │
    │ ├── Column Chunk: amount │
    │ │ └── … │
    │ └── Column Chunk: … │
    ├──────────────────────────────────────────────┤
    │ Row Group 2 │
    │ └── … (similar structure) │
    ├──────────────────────────────────────────────┤
    │ Footer │
    │ ├── Schema definition │
    │ ├── Row Group metadata │
    │ │ ├── Row count │
    │ │ ├── Column statistics (min/max/null) │
    │ │ └── Offset to column chunks │
    │ └── Key-value metadata (Iceberg: field IDs) │
    ├──────────────────────────────────────────────┤
    │ 4 bytes: Footer length │
    ├──────────────────────────────────────────────┤
    │ 4 bytes: Magic Number "PAR1" │
    └──────────────────────────────────────────────┘

    Iceberg 如何利用 Parquet 统计信息

    Iceberg 实现了四级数据跳过机制,其中最后两级依赖 Parquet 的内部结构:

    级别数据源过滤粒度是否打开文件
    1. Partition Pruning Manifest 中的分区值 整个分区(数十~数百文件)
    2. Manifest Filter Manifest 中的 lower/upper bounds 单个数据文件
    3. Row Group Filter Parquet Footer 中的 Row Group 统计 10K~100K 行 读 Footer
    4. Page Filter Parquet Page Header 中的统计 1K~10K 行 读 Page Header

    压缩算法选择指南

    压缩算法压缩比解压速度推荐场景
    ZSTD (默认) 高(3~5x) 通用推荐,2026 年最佳选择
    Snappy 中(2~3x) 极快 CPU 敏感场景
    GZIP 极高(5~8x) 存储成本极度敏感
    LZ4 低(1.5~2x) 极快 极致读取性能
    Brotli 极高(6~10x) 归档存储

    — 配置压缩算法
    ALTER TABLE demo.analytics.events SET TBLPROPERTIES (
    'write.parquet.compression-codec' = 'zstd',
    'write.parquet.compression-level' = '3' — ZSTD 压缩级别 1~22
    );


    A.3 Iceberg 与数据治理

    数据治理是企业级数据平台的核心能力。Iceberg 在规范层面提供了多项治理支持。

    数据血缘(Data Lineage)

    数据血缘追踪架构:

    ┌──────────┐ 写入 ┌──────────┐ 转换 ┌──────────┐
    │ Source │ ─────────→ │ Bronze │ ─────────→ │ Silver │
    │ MySQL │ │ orders │ Spark │ orders │
    └──────────┘ └──────────┘ └────┬─────┘

    聚合 │

    ┌────┴─────┐
    │ Gold │
    │ summary │
    └──────────┘

    Iceberg v3 Row Lineage 提供:
    ├── 行级别血缘:每行有唯一 _row_id
    ├── 时间级别血缘:_last_updated_sequence_number
    ├── 跨引擎可读:任何 Iceberg 引擎都能查询
    └── 快照级别血缘:通过 snapshots 表追踪

    配合 OpenLineage 提供:
    ├── 作业级别血缘:哪个 Spark/Flink 作业写了哪张表
    ├── Schema 变更追踪:Schema 演进历史
    └── 数据质量指标:记录数、文件大小、处理时间

    数据质量监控

    # 基于 PyIceberg 的数据质量检查框架
    from pyiceberg.catalog import load_catalog
    from datetime import datetime, timedelta

    catalog = load_catalog("demo")

    def check_table_health(table_name: str) > dict:
    """检查 Iceberg 表的健康状态"""
    table = catalog.load_table(table_name)
    snapshot = table.current_snapshot()

    health = {
    "table": table_name,
    "timestamp": datetime.now().isoformat(),
    "checks": {}
    }

    # 检查1: 数据新鲜度
    if snapshot:
    age_hours = (datetime.now().timestamp() * 1000 snapshot.timestamp_ms) / 3600000
    health["checks"]["freshness_hours"] = round(age_hours, 2)
    health["checks"]["freshness_ok"] = age_hours < 24 # 24小时内

    # 检查2: 快照数量
    num_snapshots = len(list(table.snapshots))
    health["checks"]["snapshot_count"] = num_snapshots
    health["checks"]["snapshot_ok"] = num_snapshots < 500

    # 检查3: 文件统计
    summary = snapshot.summary if snapshot else {}
    total_files = int(summary.get("total-data-files", 0))
    delete_files = int(summary.get("total-delete-files", 0))
    health["checks"]["total_data_files"] = total_files
    health["checks"]["total_delete_files"] = delete_files
    health["checks"]["file_ratio_ok"] = delete_files < total_files * 0.1

    # 检查4: Schema 一致性
    health["checks"]["schema_fields"] = len(table.schema.fields)
    health["checks"]["schema_id"] = table.schema_id

    return health

    # 执行检查
    result = check_table_health("analytics.orders")
    for check, value in result["checks"].items():
    status = "✅" if "ok" not in check or value else "❌"
    print(f" {status} {check}: {value}")

    GDPR 合规与数据删除

    — 使用 Iceberg 实现 GDPR 数据删除

    — Step 1: 识别需要删除的用户数据
    — 利用 Row Lineage 追踪受影响的数据
    SELECT _row_id, _last_updated_sequence_number, user_id
    FROM demo.analytics.user_events
    WHERE user_id = 12345;

    — Step 2: 执行删除
    DELETE FROM demo.analytics.user_events
    WHERE user_id = 12345;

    — Step 3: 验证删除
    SELECT COUNT(*) FROM demo.analytics.user_events
    WHERE user_id = 12345; — 应该返回 0

    — Step 4: 记录合规审计日志
    INSERT INTO demo.compliance.deletion_log
    VALUES (
    'GDPR_ERASURE',
    'user_12345',
    current_timestamp(),
    'operator@example.com',
    'User requested data erasure per GDPR Article 17'
    );

    — Step 5: 确保删除不可恢复(清理快照后无法 Time Travel 到删除前的状态)
    CALL demo.system.expire_snapshots(
    table => 'analytics.user_events',
    older_than => current_timestamp()
    );


    A.4 Iceberg 运维自动化脚本集

    #!/usr/bin/env python3
    """
    Iceberg 表自动化维护脚本
    适用于 Airflow / DolphinScheduler 调度
    """

    import sys
    import json
    import logging
    from datetime import datetime, timedelta
    from pyiceberg.catalog import load_catalog

    logging.basicConfig(level=logging.INFO)
    logger = logging.getLogger(__name__)

    # 配置
    CATALOG_CONFIG = {
    "type": "rest",
    "uri": "http://polaris:8181",
    "warehouse": "s3://lakehouse",
    }

    # 需要维护的表列表
    TABLES_TO_MAINTAIN = [
    "analytics.orders",
    "analytics.user_events",
    "analytics.customers",
    "ecommerce.bronze_orders",
    "ecommerce.silver_orders",
    "ecommerce.gold_daily_summary",
    ]

    def compact_table(catalog, table_name: str, target_size_mb: int = 512):
    """合并小文件"""
    logger.info(f"[Compaction] Starting for {table_name}")
    table = catalog.load_table(table_name)

    # 获取当前文件统计
    snapshot = table.current_snapshot()
    if not snapshot:
    logger.info(f"[Compaction] No snapshots found, skipping")
    return

    summary = snapshot.summary
    total_files = int(summary.get("total-data-files", 0))
    logger.info(f"[Compaction] Current data files: {total_files}")

    if total_files < 100:
    logger.info(f"[Compaction] File count below threshold, skipping")
    return

    # 执行 compaction(通过 Spark 提交)
    logger.info(f"[Compaction] Triggering rewrite for {table_name}")
    # 实际生产中这里应该通过 Spark REST API 或 Livy 提交 Spark 作业
    logger.info(f"[Compaction] Completed for {table_name}")

    def expire_snapshots(catalog, table_name: str, max_age_days: int = 7, min_snapshots: int = 10):
    """清理过期快照"""
    logger.info(f"[Expire] Starting for {table_name}")
    table = catalog.load_table(table_name)

    snapshots = list(table.snapshots)
    if len(snapshots) <= min_snapshots:
    logger.info(f"[Expire] Only {len(snapshots)} snapshots, skipping")
    return

    cutoff = datetime.now() timedelta(days=max_age_days)
    cutoff_ms = int(cutoff.timestamp() * 1000)

    expirable = [s for s in snapshots if s.timestamp_ms < cutoff_ms]
    if len(expirable) <= min_snapshots:
    logger.info(f"[Expire] Not enough expirable snapshots, skipping")
    return

    logger.info(f"[Expire] Found {len(expirable)} snapshots to expire")
    # 实际执行过期操作
    logger.info(f"[Expire] Completed for {table_name}")

    def generate_health_report(catalog) > dict:
    """生成健康报告"""
    report = {
    "timestamp": datetime.now().isoformat(),
    "tables": []
    }

    for table_name in TABLES_TO_MAINTAIN:
    try:
    table = catalog.load_table(table_name)
    snapshot = table.current_snapshot()
    summary = snapshot.summary if snapshot else {}

    table_report = {
    "name": table_name,
    "data_files": int(summary.get("total-data-files", 0)),
    "delete_files": int(summary.get("total-delete-files", 0)),
    "total_records": int(summary.get("total-records", 0)),
    "total_file_size": int(summary.get("total-file-size-bytes", 0)),
    "format_version": table.properties.get("format-version", "unknown"),
    "last_updated_ms": snapshot.timestamp_ms if snapshot else None,
    }
    report["tables"].append(table_report)

    except Exception as e:
    logger.error(f"Failed to check table {table_name}: {e}")
    report["tables"].append({
    "name": table_name,
    "error": str(e)
    })

    return report

    def main():
    """主函数"""
    logger.info("=" * 60)
    logger.info("Iceberg 表维护开始")
    logger.info("=" * 60)

    catalog = load_catalog("maintenance", **CATALOG_CONFIG)

    # Step 1: 健康检查
    logger.info("\\n— Step 1: 健康检查 —")
    report = generate_health_report(catalog)
    for t in report["tables"]:
    if "error" in t:
    logger.warning(f" ❌ {t['name']}: {t['error']}")
    else:
    logger.info(
    f" 📊 {t['name']}: "
    f"files={t['data_files']}, "
    f"deletes={t['delete_files']}, "
    f"records={t['total_records']:,}"
    )

    # Step 2: Compaction
    logger.info("\\n— Step 2: 小文件合并 —")
    for table_name in TABLES_TO_MAINTAIN:
    try:
    compact_table(catalog, table_name)
    except Exception as e:
    logger.error(f"Compaction failed for {table_name}: {e}")

    # Step 3: Snapshot 过期
    logger.info("\\n— Step 3: 快照清理 —")
    for table_name in TABLES_TO_MAINTAIN:
    try:
    expire_snapshots(catalog, table_name)
    except Exception as e:
    logger.error(f"Expire failed for {table_name}: {e}")

    logger.info("\\n" + "=" * 60)
    logger.info("Iceberg 表维护完成")
    logger.info("=" * 60)

    # 输出报告
    print(json.dumps(report, indent=2, default=str))

    if __name__ == "__main__":
    main()

    A.5 Iceberg 常见面试问题汇总

    Q1: Iceberg 如何保证 ACID 事务?

    A: Iceberg 通过原子性的元数据交换实现 ACID。每次写入产生新的 metadata.json 文件,通过 Catalog 原子地将当前快照指针从旧 metadata 切换到新 metadata。读取始终基于快照(MVCC),不会被并发写入阻塞。写入通过乐观并发控制(CAS)解决冲突。

    Q2: Iceberg 的隐藏分区与传统 Hive 分区有什么区别?

    A: Hive 分区是物理的——用户必须知道分区列并在 WHERE 子句中显式指定。Iceberg 分区是逻辑的——分区转换函数存储在元数据中,用户按原始列过滤,引擎自动推导分区裁剪。这意味着可以在线更改分区方案而不影响现有查询。

    Q3: Iceberg v3 的 Deletion Vectors 和 v2 的位置删除文件有什么区别?

    A: v2 为每次删除创建单独的 Avro 位置删除文件,列出被删除行的文件路径和行号,读取时需要合并所有删除文件,性能线性下降。v3 使用 Roaring Bitmap(存储在 Puffin 文件中)作为每个数据文件的删除向量,读取时只需检查位图,O(1) 性能,DML 性能提升约 10 倍。

    Q4: 什么场景下用 Copy-on-Write,什么场景用 Merge-on-Read?

    A: COW 适合读多写少的场景(如 BI 报表),写入时重写文件,读取时无需合并。MOR 适合写多读少的场景(如 CDC 入湖),写入时只写删除标记,读取时需要合并。v3 的 Deletion Vectors 大幅降低了 MOR 的读取开销,使得 MOR 在大多数场景下成为更好的选择。

    Q5: 如何从 v2 升级到 v3?有哪些注意事项?

    A: 升级是元数据操作,ALTER TABLE … SET TBLPROPERTIES ('format-version' = '3'),无需重写数据。注意事项:(1) 确保所有读写引擎支持 v3;(2) 旧数据文件仍然有效;(3) 新写入的文件使用 v3 特性;(4) 如有不兼容引擎,保持 v2。

    Q6: Iceberg 如何处理小文件问题?

    A: 小文件是流式写入的主要挑战。解决方案:(1) 增大 checkpoint/trigger 间隔减少文件产生频率;(2) 定期运行 Spark Actions 的 rewrite_data_files 进行 compaction;(3) 使用 BinPack/Sort/Z-Order 策略合并文件;(4) 配置自动 compaction 调度。

    Q7: REST Catalog 相比 Hive Catalog 有什么优势?

    A: REST Catalog 的优势包括:(1) HTTP 协议更通用,跨语言支持更好;(2) 支持服务端提交和冲突解决;(3) 凭证代理(Credential Vending)实现零信任安全;(4) OAuth2/JWT 原生鉴权;(5) 多表事务支持;(6) 更好的高可用扩展性。2026 年 REST Catalog 已经是事实标准。

    Q8: Iceberg 的 Schema Evolution 是如何工作的?

    A: Iceberg 使用列 ID 而非列位置来追踪 Schema 变化。每列有唯一的整数 ID,存储在元数据和 Parquet 文件中。增删改列名只需更新元数据,不需要重写数据文件。读取时,引擎通过列 ID 映射将当前 Schema 与数据文件的 Schema 对齐。类型拓宽(如 INT → BIGINT)也是安全的,因为数据值兼容。

    Q9: 什么是 Iceberg 的 Partition Evolution?为什么重要?

    A: Partition Evolution 允许在线更改分区方案而不重写数据。旧的分区规范仍然记录在元数据中,新写入使用新规范,查询引擎自动根据每个文件的分区规范进行裁剪。这意味着从日分区改为小时分区不再需要 PB 级数据重写——这是 Hive 做不到的。

    Q10: 2026 年应该选择 Iceberg 还是 Delta Lake 还是其他格式?

    A: 选择取决于场景和引擎栈:(1) 多引擎通用湖仓、生态最广 → Iceberg(2026 年默认选择);(2) 深度 Databricks 用户 → Delta Lake(UniForm 兼容 Iceberg);(3) 高频 CDC + Spark 为主 → Hudi;(4) Flink 流式优先 → Paimon。关键决策因素是你的主力计算引擎生态,不是功能列表。

    Q11: Iceberg v3 的 VARIANT 类型和直接存 JSON STRING 有什么本质区别?

    A: 本质区别在于 Shredding(碎片化存储)优化。JSON STRING 在查询时需要全量解析每个 JSON 文档,无法利用列式存储的优势。VARIANT 类型在写入时将频繁查询的字段"碎片化"为独立的类型化列存储在 Parquet 中,查询这些字段时享受与原生列相同的性能。此外,VARIANT 支持 Filter Pushdown——查询条件可以下推到文件级别跳过不相关的数据文件。简单来说:JSON STRING 是文本解析,VARIANT 是列式存储。

    Q12: 如何监控 Iceberg 表的健康状态?有哪些关键指标?

    A: 关键监控指标包括:(1) 数据文件数量——超过阈值(如 10,000)触发 compaction 告警;(2) 删除文件数量——MOR 模式下删除文件过多说明 compaction 不足;(3) 快照数量——超过阈值说明过期清理未执行;(4) 最新快照年龄——超过预期新鲜度说明写入管道可能中断;(5) 平均文件大小——远小于目标值说明小文件问题;(6) 元数据目录大小——异常增长说明清理策略需要调整。

    Q13: Iceberg 如何实现数据脱敏?

    A: 数据脱敏通常在查询层实现,而非存储层。常见方式:(1) 视图脱敏——创建视图对敏感列使用 SQL 函数脱敏(如 SUBSTR(email, 1, 3) || '***');(2) RBAC + 视图——不同角色的用户访问不同的视图;(3) Trino 行级/列级安全——在 Trino catalog 配置中设置行过滤和列掩码;(4) 应用层脱敏——在查询引擎和最终用户之间增加脱敏中间件。Iceberg 本身的 Time Travel 功能可以追溯到脱敏前的数据,因此需要确保脱敏策略与快照管理策略配合。

    Q14: REST Catalog 的凭证代理(Credential Vending)是什么?为什么重要?

    A: 凭证代理是 REST Catalog 的核心安全特性。传统方式中,计算引擎(Spark/Trino)直接持有 S3 的永久密钥(Access Key),一旦密钥泄露,攻击者可以访问整个存储桶。凭证代理的工作方式是:引擎向 REST Catalog 请求访问某张表,Catalog 验证身份和权限后,调用 AWS STS 签发临时凭证(有效期 15 分钟),该凭证的权限被 Session Policy 限定到具体表的存储前缀。引擎使用临时凭证访问 S3,15 分钟后凭证自动过期。这样即使凭证泄露,攻击窗口也只有 15 分钟,且只能访问特定表。这是零信任架构在数据湖中的标准实现。

    Q15: 从 Hive 表迁移到 Iceberg 表的最佳策略是什么?

    A: 推荐的迁移策略是渐进式迁移:(1) 第一步:创建 Iceberg 外部表——在 Iceberg 中创建指向 Hive 表数据文件的外部表(不复制数据),验证查询正确性;(2) 第二步:CTAS 迁移——使用 CREATE TABLE … AS SELECT 将 Hive 数据复制到 Iceberg 表中,同时应用更好的分区和排序策略;(3) 第三步:双写过渡期——同时写入 Hive 表和 Iceberg 表,对比查询结果;(4) 第四步:切换下游——将所有下游查询从 Hive 表切换到 Iceberg 表;(5) 第五步:停止 Hive 写入——完全切换到 Iceberg。关键注意事项:迁移时应用 Iceberg 的最佳实践(ZSTD 压缩、合理的分区策略、排序顺序),而不是简单复制 Hive 的布局。


    A.6 初学者常见问题与排错指南

    问题 1:查询 Iceberg 表时报 “Table not found” 错误

    原因通常是 Catalog 配置不正确或命名空间未创建。排查步骤:

  • 确认 Catalog URI 正确且可访问:curl http://your-catalog:8181/v1/config
  • 确认命名空间已创建:catalog.list_namespaces()
  • 确认表名格式正确:catalog.table_name(catalog.命名空间.表名)
  • 如果使用 REST Catalog,确认 warehouse 路径与对象存储匹配
  • 问题 2:写入 Iceberg 表时报 “Commit failed” 错误

    这通常是并发写入冲突。Iceberg 使用乐观并发控制,当两个 writer 同时提交时会发生冲突。解决方案:

  • 确认 commit.retry.num-workers 配置合理(推荐 2~4)
  • 减少写入并行度,增大每个 task 的数据量
  • 错开不同管道的写入时间
  • 如果冲突频繁,考虑增大分区粒度,减少不同 writer 写入同一分区的概率
  • 问题 3:查询性能突然下降

    最常见的原因是小文件累积。排查步骤:

  • 检查数据文件数量:SELECT * FROM table.snapshots 查看 total-data-files
  • 如果文件数 > 10,000,执行 compaction
  • 检查快照数量:如果快照数 > 500,执行 expire_snapshots
  • 检查是否存在大量删除文件(MOR 模式)
  • 确认查询条件命中了分区键(分区裁剪是否生效)
  • 问题 4:Flink 写入 Iceberg 后查询不到数据

    Flink 使用 checkpoint 机制提交数据到 Iceberg。如果 checkpoint 未完成,数据对外不可见。排查步骤:

  • 确认 Flink 作业的 checkpoint 已启用且正常完成
  • 检查 Flink Web UI 中的 checkpoint 状态
  • 确认 Flink 的 iceberg catalog 配置与查询引擎一致
  • 确认 Flink checkpoint 间隔设置合理(推荐 30~60 秒)
  • 问题 5:Time Travel 查询返回空结果

    可能原因:(1) 目标快照已过期被清理;(2) 时间点没有对应的快照。解决方案:

  • 查看可用快照:SELECT * FROM table.snapshots ORDER BY committed_at
  • 调整快照保留策略:'history.expire.max-snapshot-age-ms' = '604800000'(7天)
  • 使用精确的快照 ID 而非时间戳进行查询
  • 问题 6:ALTER TABLE 操作后查询报错 “Schema mismatch”

    可能是 Schema 演进后旧数据文件与新 Schema 不兼容。Iceberg 通过列 ID 映射解决大部分 Schema 差异,但某些类型拓宽操作可能需要额外处理。确保:

  • 只做兼容的类型拓宽(INT→BIGINT, FLOAT→DOUBLE)
  • 不支持的类型转换(STRING↔INT)需要先 ALTER COLUMN TYPE 再重写数据
  • 删除列后查询旧快照时,被删除的列会显示为 NULL(这是正常行为)
  • 初学者学习路径建议

    推荐学习路径(4 周计划):

    第 1 周:基础概念
    ├── 阅读本指南第 1~3 章
    ├── Docker 搭建本地环境
    ├── 完成基本的 CREATE/INSERT/SELECT 操作
    └── 理解五层元数据结构

    第 2 周:SQL 操作
    ├── 阅读本指南第 5~7 章
    ├── 练习 DDL/DML 操作
    ├── 尝试 Time Travel 和分支操作
    └── 学习 Schema Evolution

    第 3 周:Python SDK
    ├── 阅读本指南第 8 章
    ├── 使用 PyIceberg 完成所有操作
    ├── 尝试连接不同的 Catalog
    └── 阅读 Iceberg 官方文档

    第 4 周:进阶实践
    ├── 阅读本指南第 9~16 章
    ├── 搭建 Flink + Iceberg 流式管道
    ├── 实现 CDC 入湖方案
    ├── 练习性能调优和维护操作
    └── 阅读 Iceberg Summit 2026 技术分享

    A.7 Iceberg 核心术语速查表

    术语英文含义
    表格式 Table Format 定义数据文件如何组织为一致、可查询的表的元数据规范
    快照 Snapshot 表在某一时刻的完整视图,是 Time Travel 的基础
    清单列表 Manifest List Avro 文件,列出某个快照包含的所有 Manifest File
    清单文件 Manifest File Avro 文件,列出数据文件及其分区值和列统计信息
    数据文件 Data File Parquet/ORC 格式的实际数据存储文件
    元数据文件 Metadata File JSON 文件,包含表的完整状态(schema、分区、快照等)
    隐藏分区 Hidden Partitioning 分区转换存储在元数据中,用户无需感知分区列
    分区演进 Partition Evolution 在线更改分区方案而不重写数据文件
    Schema 演进 Schema Evolution 增删改列名而无需重写数据文件,基于列 ID 映射
    时间旅行 Time Travel 查询表的历史快照数据
    写时复制 Copy-on-Write (COW) 删除/更新时重写整个数据文件,读取无需合并
    读时合并 Merge-on-Read (MOR) 删除/更新时写入删除标记,读取时合并
    删除向量 Deletion Vector (v3) Roaring Bitmap 位图标记被删除的行,替代位置删除文件
    行血缘 Row Lineage (v3) 每行的唯一 ID 和最后修改序列号
    凭证代理 Credential Vending REST Catalog 为引擎签发临时存储凭证
    文件重写 Compaction / Rewrite 合并小文件以提升查询性能
    快照过期 Expire Snapshots 清理超过保留期限的历史快照
    孤立文件 Orphan Files 存在于存储但不在任何快照中引用的文件
    分支 Branch 基于某个快照的可变引用,支持独立开发
    标签 Tag 基于某个快照的不可变引用,用于标记重要版本
    表格式版本 Format Version 磁盘上的规范版本(v1/v2/v3),区别于库版本
    元数据目录 Metadata Directory 存储 metadata.json、manifest list、manifest file 的目录
    Puffin 文件 Puffin File 存储二进制统计信息的文件格式,v3 中用于存储 Deletion Vector
    文件规划 File Planning 查询引擎根据元数据统计信息决定需要读取哪些数据文件的过程
    服务端提交 Server-side Commit REST Catalog 在服务器端处理提交冲突,减少客户端重试
    乐观并发控制 Optimistic Concurrency Control 基于 CAS 的并发写入控制策略,无需加锁

    A.8 常用命令速查表

    # Spark SQL 常用命令
    CREATE TABLE ... USING iceberg TBLPROPERTIES ('format-version' = '3');
    ALTER TABLE ... ADD COLUMNS (...);
    ALTER TABLE ... SET TBLPROPERTIES ('key' = 'value');
    INSERT INTO ... SELECT ...;
    INSERT OVERWRITE ... PARTITION (...);
    MERGE INTO ... USING ... ON ... WHEN MATCHED ...;
    DELETE FROM ... WHERE ...;
    UPDATE ... SET ... WHERE ...;
    SELECT * FROM ... VERSION AS OF <snapshot_id>;
    SELECT * FROM ... TIMESTAMP AS OF '<timestamp>';
    SELECT * FROM ... TABLE_NAME.snapshots;
    SELECT * FROM ... TABLE_NAME.history;
    CALL catalog.system.rewrite_data_files(table => '…');
    CALL catalog.system.expire_snapshots(table => '…');
    CALL catalog.system.remove_orphan_files(table => '…');

    B. PyIceberg 常用 API 速查

    # 连接
    catalog = load_catalog("name", type="rest", uri="http://…")

    # 命名空间
    catalog.create_namespace("ns")
    catalog.list_namespaces()

    # 表操作
    catalog.create_table("ns.table", schema=schema, properties={...})
    catalog.load_table("ns.table")
    catalog.list_tables("ns")
    catalog.drop_table("ns.table")

    # 读写
    table.append(data) # PyArrow Table
    table.overwrite(data)
    table.scan().to_arrow()
    table.scan(row_filter="…").to_pandas()

    # Schema 演进
    with table.update_schema() as update:
    update.add_column("col", StringType())
    update.rename_column("old", "new")
    update.delete_column("col")

    # 属性
    with table.update_properties() as update:
    update.set("key", "value")

    C. 参考资料

    • Apache Iceberg 官方文档
    • Apache Iceberg 1.11.0 Release Notes
    • Apache Iceberg 规范(V3)
    • Apache Polaris 官方文档
    • PyIceberg 文档
    • The Apache Iceberg Market in the Middle of 2026

    赞(0)
    未经允许不得转载:171主机测评 » Apache Iceberg 学习指南
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址