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 与 Delta Lake(delta-rs)如何选择?
如何安装和配置 PyIceberg?
如何用 PyIceberg 创建第一张 Iceberg 表?
如何写入并读取 Iceberg 表数据?
PyIceberg 如何进行 Schema 演化?
如何做分区演化(Partition Evolution)?
Iceberg 时间旅行与快照管理
与 DuckDB、Polars、Pandas 集成
生产环境:AWS Glue REST Catalog 与 S3 Tables
数据管道测试:让 PyIceberg 不再半夜报警
性能实测:PyIceberg vs Spark on Iceberg
上生产前的检查清单
常见问题解答
什么是 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 Iceberg Delta 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、Doris Spark、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"
警告: 不要在生产环境把 PyIceberg 和 Iceberg Java 库版本混用时忽略兼容性矩阵。PyIceberg 0.11.x 对应 Iceberg 表规范 v2,如果你写入的表要被老版本 Spark 3.3 及以下读取,请务必测试。
配置 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 文档 里对所有变换(identity、truncate、bucket、year/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())
提示: 永远优先用 selected_fields 做列裁剪,尤其是宽表。Iceberg 的 Parquet 文件以列存储,只读取需要的列可以让 IO 下降一个数量级。这不是玄学优化。在我们线上一张 200 列的事件表上,从"读全部"改成"读 6 列"后,同一个查询 S3 出口流量从 8 GB 降到 380 MB。
PyIceberg 如何进行 Schema 演化?
Schema 演化是 Iceberg 相对 Hive 表的核心优势之一。因为每个字段有稳定 ID,你可以安全地新增列、删除列、改名、放宽类型(比如 int 到 long),历史 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 默认只允许"非破坏性"变更。如果你尝试删列或收窄类型(比如 long 到 int),它会直接抛错。你需要显式传 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")
注意: 分区演化后,跨越旧或新规则边界的查询会命中两种物理布局。查询引擎(DuckDB、Trino)能正确合并结果,但计划复杂度上升。如果 60% 以上的数据仍在旧规则下,建议先跑一次 rewrite_data_files 后台任务再切换。
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_null、unique、relationships 语义。
警告: 回填(backfill)是 PyIceberg 生产里最容易出事的操作。默认的 overwrite 会替换整个匹配分区,如果你的谓词写错,可能会误删几个月的数据。回填必须做三件事:先在 dry-run 分支验证;写入前拍一次快照(表本身就是快照,但要记录 ID);用 overwrite(overwrite_filter=...) 精确指定范围。
性能实测:PyIceberg vs Spark on Iceberg
2026 年一份被广泛引用的对比(apache/iceberg-python 仓库 里也有类似基准)显示,在 1 亿行 Parquet 数据上做谓词下推加列裁剪查询,PyIceberg 加 DuckDB 的冷启动到出结果耗时约 3.2 秒,Spark on Iceberg(单机 local[*])大约 47 秒。当然 Spark 在超过百亿行、复杂 join 场景下仍有优势,但对绝大多数部门级分析和 ML 特征工程,PyIceberg 已经完全够用而且成本低一个数量级。
上生产前的检查清单
Catalog 后端 :本地用 SqlCatalog 测试,生产用 REST(Polaris/Nessie/Glue REST),避免多写者并发提交冲突。
写入并发 :设置 commit.retry.num-retries(默认 4,繁忙表建议 8 到 10)。
文件大小 :write.target-file-size-bytes 设 128 MB 到 256 MB,太小会造成 S3 GET 风暴,太大会拖慢 scan。
压缩 :ZSTD level 3 是 2026 年默认最佳权衡。
小文件合并 :rewrite_data_files 每天跑一次;expire_snapshots 保留 7 天;remove_orphan_files 每周跑一次。
可观测性 :把每次 append 后的 snapshot_id、行数、字节数打点到 Prometheus 或 DataDog。
数据契约测试 :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 处理指标和维度建模,职责很清晰。