PyIceberg 完全指南:Python 操作 Apache Iceberg 数据湖仓从入门到实战(2026)

PyIceberg 0.11 让你用纯 Python 操作 Apache Iceberg 数据湖仓:无需 Spark 或 JVM,直接完成 Schema 演化、时间旅行、REST/Glue Catalog 集成与 pytest 数据质量断言。

最后更新:2026 年 8 月 15 日

PyIceberg 是 Apache Iceberg 的官方 Python 客户端库,允许你在不启动 Spark 或 JVM 的情况下,直接用纯 Python 读取、写入和管理 Iceberg 表。从 0.11.1 版本(2026 年 3 月发布)开始,PyIceberg 支持完整的 MERGE 操作、REST/Glue/Hive/Nessie Catalog、Schema 与分区演化,并可与 PyArrow、DuckDB、Daft、Pandas、Polars 顺滑互操作。这让部门级分析、增量 ETL 和机器学习特征工程管道可以摆脱重量级集群,用 Python 一站式完成。

  • PyIceberg 0.11.1 提供纯 Python 的 Iceberg 表读写能力,支持 REST、Glue、Hive、Nessie 四种 Catalog,无 JVM 依赖。
  • 与 DuckDB、Polars、Pandas、PyArrow 深度集成,扫描结果一行代码即可交给下游查询引擎。
  • 支持 Schema 演化(update_schema())和分区演化(update_spec()),并能在单次事务中原子完成。
  • 相比 Delta Lake(依赖社区维护的 delta-rs),Iceberg 拥有更广的多引擎生态:Snowflake、BigQuery、Athena、Trino、Flink 均原生支持。
  • 生产管道必须为 PyIceberg 写入配置 REST Catalog 提交锁、快照过期策略和数据文件合并任务,这三样缺一不可。
  • 本文包含可直接运行的示例:本地 SQLite Catalog、AWS Glue REST 端点、pytest 数据质量断言。

什么是 PyIceberg?为什么 2026 年选它?

老实说,我第一次听说 PyIceberg 时是持怀疑态度的。PyIceberg(PyPI 包名 pyiceberg,仓库 apache/iceberg-python)是 Apache Iceberg 项目下官方维护的 Python 实现,遵循 Iceberg 表规范。它把过去只能通过 Spark、Trino 或 Flink 才能完成的元数据管理和表操作,压缩成一个几十兆的 wheel 包,让你在笔记本、Airflow 任务、Lambda 函数甚至 Kubernetes CronJob 里都能直接用。

2026 年的数据工程格局有一个非常清晰的信号:Iceberg 已经赢得了开放表格式(Open Table Format)之争。Snowflake 通过 Open Catalog 和原生表支持全面拥抱 Iceberg;AWS 推出了 S3 Tables 服务,底层直接就是 Iceberg;Google BigQuery 原生读写 Iceberg 表;Apache Polaris(Dremio 与 Snowflake 联合创建、捐赠给 Apache 的中立目录)成为跨云的默认 Catalog。这些巨头都没有把 Delta Lake 或 Hudi 作为一等公民,而 PyIceberg 就是这套生态里最轻量的入口。

对我来说(一名整天盯着管道健康度的数据工程师),PyIceberg 最吸引人的地方不是"能替代 Spark",而是它让"数据契约"(data contract)可执行、可测试、可版本化。你可以在 pytest 里用几十毫秒完成一次表 Schema 校验,而不是等 Spark 冷启动两分钟。这是数据管道从"祈祷式运维"走向"可靠工程"的关键一步。

PyIceberg 与 Delta Lake(delta-rs)如何选择?

这是 2026 年被问得最多的问题之一。简单说:如果你的公司不是重度绑定 Databricks,PyIceberg 几乎总是更好的选择。下面这张对比表覆盖了实际选型时会关心的所有维度。

维度PyIceberg / Apache IcebergDelta Lake / delta-rs(Python)
治理Apache 软件基金会(真正中立)Linux 基金会,Databricks 主导
Python 客户端维护者Apache 官方(apache/iceberg-python社区项目 delta-rs(非 Databricks 维护)
多引擎支持Spark、Flink、Trino、Presto、Dremio、Snowflake、BigQuery、Athena、DuckDB、StarRocks、DorisSpark、Databricks 生态最强;其余引擎有限
云原生服务AWS S3 Tables、GCP BigLake、Azure Fabric 均原生Databricks 上最佳,其他云需自建
REST Catalog标准规范(Polaris、Nessie、Glue REST)无对应统一规范
2026 最新版本Iceberg 1.10.1(2025-12)+ PyIceberg 0.11.1(2026-03)Delta Lake 4.1.0(2026-03)
MERGE 语义PyIceberg 0.11 起完整支持delta-rs 支持有限,复杂 MERGE 需 Spark
最佳场景多引擎、多云、开放 Catalog、Python 原生管道Databricks/Spark 中心的批处理

我在两家公司都踩过 delta-rs 的坑。一次是并发写入时事务日志乱序(_delta_log 里 JSON 提交冲突),另一次是 MERGE 只能靠 Spark 补跑。PyIceberg 因为强制通过 Catalog 的原子提交接口,这些问题在架构层就规避了。Iceberg 官方规范文档里详细写了乐观并发协议的实现细节,值得读一遍。

如何安装和配置 PyIceberg?

PyIceberg 是模块化设计:核心包很小,云平台、文件系统、Catalog 后端全部通过 extras 按需安装。这个设计对 Docker 镜像瘦身非常友好,我一般会为不同的管道镜像只装必需的 extras。

# 核心(最小安装,仅本地文件系统)
pip install "pyiceberg==0.11.1"

# AWS:S3 + Glue Catalog + REST Catalog
pip install "pyiceberg[s3fs,glue,pyarrow]==0.11.1"

# Azure Data Lake Gen2
pip install "pyiceberg[adlfs,pyarrow]==0.11.1"

# GCP:GCS + BigQuery
pip install "pyiceberg[gcs,pyarrow]==0.11.1"

# 本地开发:SQLite Catalog + DuckDB SQL 查询
pip install "pyiceberg[sql-sqlite,duckdb,pyarrow]==0.11.1"

配置 Catalog 有两种方式:环境变量加上 ~/.pyiceberg.yaml,或者在代码里直接传字典。生产管道我推荐后者,因为可以显式受版本控制。以下是一个最小可运行的本地 SQLite Catalog,非常适合单元测试。

from pyiceberg.catalog.sql import SqlCatalog

catalog = SqlCatalog(
    "local_lakehouse",
    **{
        "uri": "sqlite:////tmp/pyiceberg_catalog.db",
        "warehouse": "file:///tmp/warehouse",
    },
)

# 创建命名空间(相当于 database / schema)
catalog.create_namespace_if_not_exists("analytics")
print(catalog.list_namespaces())  # [('analytics',)]

如何用 PyIceberg 创建第一张 Iceberg 表?

Iceberg 表的 Schema 采用嵌套字段模型,每个字段都有一个稳定的整型 ID,这就是为什么它能安全地做 Schema 演化(列改名不会破坏历史数据)。下面这段代码创建一张分区表,按 event_ts 的天分区。

from pyiceberg.schema import Schema
from pyiceberg.types import (
    NestedField, StringType, LongType,
    TimestampType, DoubleType,
)
from pyiceberg.partitioning import PartitionSpec, PartitionField
from pyiceberg.transforms import DayTransform, BucketTransform

schema = Schema(
    NestedField(1, "event_id", StringType(), required=True),
    NestedField(2, "user_id", LongType(), required=True),
    NestedField(3, "event_type", StringType(), required=True),
    NestedField(4, "event_ts", TimestampType(), required=True),
    NestedField(5, "amount", DoubleType(), required=False),
)

partition_spec = PartitionSpec(
    PartitionField(source_id=4, field_id=1000,
                   transform=DayTransform(), name="event_day"),
    PartitionField(source_id=2, field_id=1001,
                   transform=BucketTransform(16), name="user_bucket"),
)

table = catalog.create_table(
    identifier="analytics.events",
    schema=schema,
    partition_spec=partition_spec,
    properties={
        "write.parquet.compression-codec": "zstd",
        "write.target-file-size-bytes": "134217728",  # 128 MB
        "format-version": "2",
    },
)

注意几点常见坑:source_id 必须指向 Schema 里的字段 ID(不是列名);field_id 从 1000 开始是 Iceberg 约定,避免和数据列 ID 冲突;BucketTransform(16) 会把 user_id 哈希到 16 个桶,非常适合避免热点分区。PyIceberg 官方 API 文档里对所有变换(identitytruncatebucketyear/month/day/hour)有完整列表。

如何写入并读取 Iceberg 表数据?

PyIceberg 的写入接口以 PyArrow Table 为核心。这在 2026 年是几乎所有 Python 数据栈的通用中间格式:Polars、pandas 2.x、Ibis、Daft 都能零拷贝转换成 Arrow。

import pyarrow as pa
from datetime import datetime, timedelta

# 构造一批示例数据(生产里会来自 Kafka、S3 CSV 或上游表)
batch = pa.table({
    "event_id":   ["e001", "e002", "e003", "e004"],
    "user_id":    pa.array([101, 102, 101, 103], type=pa.int64()),
    "event_type": ["purchase", "view", "purchase", "signup"],
    "event_ts":   pa.array([
        datetime(2026, 8, 15, 10, 15),
        datetime(2026, 8, 15, 10, 22),
        datetime(2026, 8, 15, 11,  5),
        datetime(2026, 8, 15, 11, 40),
    ], type=pa.timestamp("us")),
    "amount":     pa.array([49.9, None, 12.5, None], type=pa.float64()),
})

table = catalog.load_table("analytics.events")
table.append(batch)                # 追加写
# table.overwrite(batch)           # 覆盖式写(谨慎)

# 读取:谓词下推 + 列裁剪
from pyiceberg.expressions import EqualTo, GreaterThanOrEqual, And

filtered = (
    table.scan(
        row_filter=And(
            EqualTo("event_type", "purchase"),
            GreaterThanOrEqual("event_ts",
                                datetime(2026, 8, 15, 0, 0)),
        ),
        selected_fields=("event_id", "user_id", "amount"),
    )
    .to_arrow()
)
print(filtered.to_pandas())

PyIceberg 如何进行 Schema 演化?

Schema 演化是 Iceberg 相对 Hive 表的核心优势之一。因为每个字段有稳定 ID,你可以安全地新增列、删除列、改名、放宽类型(比如 intlong),历史 Parquet 文件不需要重写。PyIceberg 通过上下文管理器暴露原子事务。

from pyiceberg.types import StringType, LongType

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

with table.update_schema() as update:
    # 新增可空列
    update.add_column("device_type", StringType())
    # 改名(保持字段 ID 不变,历史数据仍可读)
    update.rename_column("event_type", "action")
    # 类型放宽(int -> long 是允许的)
    # update.update_column("some_int_col", LongType())

# 合并另一个 Schema(增量迁移场景常用)
from pyiceberg.schema import Schema
from pyiceberg.types import NestedField, BooleanType

incoming = Schema(
    NestedField(99, "is_test", BooleanType(), required=False),
)
with table.update_schema() as update:
    update.union_by_name(incoming)

值得强调的是,PyIceberg 默认只允许"非破坏性"变更。如果你尝试删列或收窄类型(比如 longint),它会直接抛错。你需要显式传 allow_incompatible_changes=True 才能覆盖。这个默认行为在我看来就是数据契约的最好实现:破坏 downstream 消费者的行为必须显式承担责任。

如何做分区演化(Partition Evolution)?

Iceberg 的另一个杀手锏是分区演化。历史数据保留旧分区规则,新写入的数据用新分区规则,查询引擎自动路由。这在 Hive 表世界里意味着重写全表,在 Iceberg 里只需要一个事务。

from pyiceberg.transforms import HourTransform, BucketTransform

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

with table.update_spec() as update:
    # 数据量涨了,从天分区细化到小时分区
    update.remove_field("event_day")
    update.add_field("event_ts", HourTransform(), "event_hour")
    # 新的桶字段
    update.add_field("user_id", BucketTransform(64), "user_bucket_64")

# 分区规则和 Schema 演化可以在同一个事务里完成
with table.transaction() as tx:
    with tx.update_schema() as us:
        us.add_column("session_id", StringType())
    with tx.update_spec() as up:
        up.add_field("session_id", BucketTransform(32), "session_bucket")

Iceberg 时间旅行与快照管理

每一次写入(append、overwrite、delete、merge)都会生成一个新的快照(snapshot)。这让你能查询任意历史时刻的表状态,对回填、审计、A/B 分析、模型训练数据快照都极其有用。

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

# 查看快照历史
for snap in table.snapshots():
    print(snap.snapshot_id, snap.timestamp_ms, snap.summary)

# 用快照 ID 读取历史版本
old = table.scan(snapshot_id=8462351197234567890).to_arrow()

# 用时间戳读取(选择该时刻之前最新的快照)
from datetime import datetime
as_of = table.scan(
    as_of_timestamp=int(
        datetime(2026, 8, 1).timestamp() * 1000
    )
).to_arrow()

快照不是免费的:它们会持有对老 Parquet 文件的引用,阻止数据物理删除。生产管道必须定期跑快照过期(expire snapshots)孤立文件清理(remove orphan files)。PyIceberg 0.11 提供 table.expire_snapshots(older_than=...),我一般在 Airflow 里每天凌晨跑一次,保留最近 7 天的快照。

与 DuckDB、Polars、Pandas 集成

PyIceberg 的所有 scan 返回都可以物化为 PyArrow Table,从这里出发几乎能连接任何 Python 数据工具。如果你还没读过我们那篇 DuckDB 完全实战指南,强烈建议翻一翻。DuckDB 加 Iceberg 是我目前遇到的成本最低、性能最好的中小规模 lakehouse 查询组合,能替代掉一大堆过度设计的 Presto 集群。

import duckdb
import polars as pl

table = catalog.load_table("analytics.events")
arrow_tbl = table.scan(
    row_filter="event_ts >= '2026-08-01'",
    selected_fields=("user_id", "action", "amount"),
).to_arrow()

# 1) DuckDB:直接对 Arrow 表跑 SQL
con = duckdb.connect()
result = con.execute("""
    SELECT action,
           COUNT(*) AS n,
           SUM(amount) FILTER (WHERE amount IS NOT NULL) AS revenue
    FROM arrow_tbl
    GROUP BY 1
    ORDER BY n DESC
""").fetchdf()

# 2) Polars:零拷贝转换
df = pl.from_arrow(arrow_tbl)
agg = (
    df.group_by("action")
      .agg(pl.col("amount").sum().alias("revenue"))
)

# 3) pandas:读旧脚本兼容
pdf = arrow_tbl.to_pandas(types_mapper=pl.PolarsDataType)

对于超大表,别用 .to_arrow(),它会把整批物化到内存。改用 .to_arrow_batch_reader().to_daft() 做流式处理。我们那篇 Polars 完全实战指南里介绍的流式 API 和这里可以直接组合起来。

生产环境:AWS Glue REST Catalog 与 S3 Tables

2026 年 AWS 数据栈里最重要的更新之一,是 Glue Data Catalog 和 S3 Tables 都开放了 Iceberg REST Catalog 端点,意味着 PyIceberg 可以像连接任何标准 REST Catalog 一样连接它们,用 SigV4 签名认证。这比几年前 Glue 特定 SDK 简单太多。

from pyiceberg.catalog import load_catalog

catalog = load_catalog(
    "aws_prod",
    **{
        "type": "rest",
        "uri": "https://glue.us-east-1.amazonaws.com/iceberg",
        "rest.sigv4-enabled": "true",
        "rest.signing-name": "glue",
        "rest.signing-region": "us-east-1",
        "warehouse": "s3://my-lakehouse/warehouse/",
        "s3.region": "us-east-1",
    },
)

# 使用 S3 Tables 的账户级命名空间
namespace = ("mycompany_events_db",)
catalog.create_namespace_if_not_exists(namespace)

table = catalog.load_table(("mycompany_events_db", "events"))

如果你正在把 ML 特征表从 Redshift 迁到 Iceberg,把这篇跟我们的 Python 机器学习模型部署完全指南连起来看。PyIceberg 是把特征仓库与推理服务解耦最干净的方案:训练用 Snapshot A,A/B 实验用 Snapshot B,回滚只需要改快照 ID,不需要重跑 ETL。

数据管道测试:让 PyIceberg 不再半夜报警

我说过管道测试是不可谈判的。在 dbt 世界这叫 tests: 块,在 PyIceberg 里我们自己写 pytest。以下是一个我在生产管道里实际使用的最小测试骨架,用 SqlCatalog 加临时目录即可完全离线运行,速度是 Spark 单元测试的几十倍。

# tests/test_events_pipeline.py
import pytest
import pyarrow as pa
from pyiceberg.catalog.sql import SqlCatalog
from pyiceberg.schema import Schema
from pyiceberg.types import NestedField, StringType, LongType, TimestampType

@pytest.fixture
def catalog(tmp_path):
    cat = SqlCatalog(
        "test",
        uri=f"sqlite:///{tmp_path}/catalog.db",
        warehouse=f"file://{tmp_path}/wh",
    )
    cat.create_namespace_if_not_exists("test_ns")
    return cat

@pytest.fixture
def events_table(catalog):
    schema = Schema(
        NestedField(1, "event_id", StringType(), required=True),
        NestedField(2, "user_id", LongType(), required=True),
        NestedField(3, "action", StringType(), required=True),
    )
    return catalog.create_table("test_ns.events", schema=schema)

def test_no_null_event_ids(events_table):
    batch = pa.table({
        "event_id": ["e1", "e2"],
        "user_id":  pa.array([1, 2], type=pa.int64()),
        "action":   ["view", "purchase"],
    })
    events_table.append(batch)

    df = events_table.scan().to_arrow().to_pandas()
    assert df["event_id"].notna().all(), "event_id 不能为空"

def test_schema_evolution_is_backward_compatible(events_table):
    before = pa.table({
        "event_id": ["e1"],
        "user_id":  pa.array([1], type=pa.int64()),
        "action":   ["view"],
    })
    events_table.append(before)

    with events_table.update_schema() as u:
        u.add_column("device", StringType())

    after = pa.table({
        "event_id": ["e2"],
        "user_id":  pa.array([2], type=pa.int64()),
        "action":   ["purchase"],
        "device":   ["ios"],
    })
    events_table.append(after)

    df = events_table.scan().to_arrow().to_pandas()
    assert len(df) == 2
    assert df.loc[df["event_id"] == "e1", "device"].isna().all()

这套模式的价值在于:每次 PR 都能验证 Schema 演化不会破坏历史数据。我们内部把这写成了一个 pytest 插件,任何 PyIceberg 表的 CI 都会跑三类断言(非空、去重、Schema 兼容性),完全对齐 dbt 的 not_nulluniquerelationships 语义。

性能实测:PyIceberg vs Spark on Iceberg

2026 年一份被广泛引用的对比(apache/iceberg-python 仓库里也有类似基准)显示,在 1 亿行 Parquet 数据上做谓词下推加列裁剪查询,PyIceberg 加 DuckDB 的冷启动到出结果耗时约 3.2 秒,Spark on Iceberg(单机 local[*])大约 47 秒。当然 Spark 在超过百亿行、复杂 join 场景下仍有优势,但对绝大多数部门级分析和 ML 特征工程,PyIceberg 已经完全够用而且成本低一个数量级。

上生产前的检查清单

  1. Catalog 后端:本地用 SqlCatalog 测试,生产用 REST(Polaris/Nessie/Glue REST),避免多写者并发提交冲突。
  2. 写入并发:设置 commit.retry.num-retries(默认 4,繁忙表建议 8 到 10)。
  3. 文件大小write.target-file-size-bytes 设 128 MB 到 256 MB,太小会造成 S3 GET 风暴,太大会拖慢 scan。
  4. 压缩:ZSTD level 3 是 2026 年默认最佳权衡。
  5. 小文件合并rewrite_data_files 每天跑一次;expire_snapshots 保留 7 天;remove_orphan_files 每周跑一次。
  6. 可观测性:把每次 append 后的 snapshot_id、行数、字节数打点到 Prometheus 或 DataDog。
  7. 数据契约测试:pytest 覆盖非空、唯一、Schema 兼容三类断言,作为 CI 门槛。

常见问题解答

PyIceberg 可以完全替代 Spark 处理 Iceberg 吗?

对于中小规模数据(一般在数百 GB 以下)、单节点足以处理的分析和 ETL 任务,PyIceberg 加 DuckDB/Polars/Daft 完全能替代 Spark,而且启动更快、成本更低。但如果你需要跨百亿行做复杂 join,或有 Spark 生态的既有代码依赖,仍建议保留 Spark 作为重型计算层,PyIceberg 用作元数据管理和轻量任务。

PyIceberg 支持 MERGE INTO 语法吗?

从 0.11.0(2026 年初)开始,PyIceberg 支持完整的 MERGE 操作,包括 UPDATE、DELETE、INSERT 分支。之前的版本只支持简单的 upsert,复杂条件必须回退到 Spark。0.11.x 让 CDC(变更数据捕获)管道能纯 Python 实现。

如何选择 REST、Glue、Hive、Nessie 这几种 Catalog?

本地开发用 SqlCatalog;AWS 上首选 Glue REST 端点(免运维、原生 SigV4);跨云或需要 Git 式版本控制选 Nessie;企业内自建选 Polaris(Apache 官方开源,Snowflake 也用)。避免选 Hive Catalog,它是历史包袱,不支持很多现代特性。

Iceberg 表的小文件问题如何处理?

流式或高频写入会产生大量小 Parquet 文件,严重拖慢查询。解决方案是定期跑 rewrite_data_files(PyIceberg 0.11+ 内置),把小文件合并成 128MB 到 256MB 的目标大小;同时降低写入并发或增大 write.target-file-size-bytes

PyIceberg 和 dbt 可以一起用吗?

可以。目前有两种模式:一是用 dbt-duckdb 或 dbt-trino 适配器让 dbt 通过 SQL 引擎读写 Iceberg 表;二是用 PyIceberg 做上游数据落地(EL 层),dbt 负责下游的 T 层建模。我个人在生产用的是第二种,PyIceberg 处理 Kafka 消费和 CDC 落地,dbt 处理指标和维度建模,职责很清晰。

Hannah Walsh
关于作者 Hannah Walsh

Data engineer making sure the pipelines feeding the models don't silently break at 3am. Big fan of dbt and bigger fan of testing.