基于Pyspark對(duì)Apache Iceberg核心功能的使用實(shí)踐指南
本文系統(tǒng)測(cè)試了 Apache Iceberg 表格式的核心功能。實(shí)驗(yàn)環(huán)境為本地 PySpark + Hadoop Catalog。
Apache Iceberg 是一種開(kāi)放表格式,為數(shù)據(jù)湖提供數(shù)據(jù)庫(kù)級(jí)別的 ACID 事務(wù)和高級(jí)功能。不同于傳統(tǒng)的 Hive 表,Iceberg 將元數(shù)據(jù)與數(shù)據(jù)分離,實(shí)現(xiàn)了隱藏分區(qū)、時(shí)間旅行、Schema 演進(jìn)等能力。
| 特性 | 說(shuō)明 |
|---|---|
| ACID 事務(wù) | 寫(xiě)入操作原子性,讀取快照隔離 |
| 隱藏分區(qū) | 分區(qū)對(duì)用戶(hù)透明,查詢(xún)規(guī)劃器自動(dòng)處理分區(qū)裁剪 |
| Schema 演進(jìn) | 支持添加、刪除、重命名、 reorder 列,零數(shù)據(jù)重寫(xiě) |
| 分區(qū)演進(jìn) | 可更改已有表的分區(qū)策略,無(wú)需遷移數(shù)據(jù) |
| 時(shí)間旅行 | 基于快照的歷史版本查詢(xún),支持 VERSION AS OF 和 TIMESTAMP AS OF |
| 行級(jí)操作 | 支持 UPDATE、DELETE、MERGE INTO |
| 開(kāi)放格式 | 支持 Parquet、Avro、ORC 等主流列式存儲(chǔ)格式 |
核心術(shù)語(yǔ)
| 術(shù)語(yǔ) | 說(shuō)明 |
|---|---|
| Schema | 表的字段定義(列名、類(lèi)型) |
| Partition Spec | 分區(qū)規(guī)范,定義如何從數(shù)據(jù)字段派生分區(qū)值 |
| Snapshot | 表在某一時(shí)刻的完整狀態(tài)快照 |
| Manifest List | 清單列表文件,記錄屬于某快照的所有 manifest 文件 |
| Manifest | 清單文件,記錄該快照包含的所有數(shù)據(jù)文件和刪除文件 |
| Data File | 實(shí)際存儲(chǔ)表數(shù)據(jù)的文件(Parquet/Avro/ORC) |
| Delete File | 記錄被刪除行的文件,用于 Merge-on-Read |
| Metadata File | 元數(shù)據(jù) JSON 文件,記錄表結(jié)構(gòu)、分區(qū)規(guī)范、快照列表 |
選型說(shuō)明
- PySpark 3.5.6:PySpark 與 Iceberg 集成最成熟,支持完整的 Iceberg SQL 語(yǔ)法和 DataFrame API。
- Hadoop Catalog:使用本地文件系統(tǒng)作為元數(shù)據(jù)存儲(chǔ),生產(chǎn)環(huán)境可替換為 Hive Metastore 或 AWS Glue。
- Iceberg 1.8.1:與 Spark 3.5.x 兼容良好。
SparkSession 配置詳解
spark.jars.packages:引入 Iceberg Spark 運(yùn)行時(shí) JAR,自動(dòng)下載依賴(lài)spark.sql.extensions:注冊(cè) Iceberg SQL 擴(kuò)展,支持 Iceberg 專(zhuān)用語(yǔ)法spark.sql.catalog.{CATALOG}:配置 Iceberg Catalog 實(shí)現(xiàn)類(lèi)spark.sql.catalog.{CATALOG}.type:指定 Catalog 類(lèi)型為 Hadoop(文件系統(tǒng)后端)spark.sql.catalog.{CATALOG}.warehouse:指定本地倉(cāng)庫(kù)路徑/tmp/iceberg-warehousespark.sql.session.timeZone:統(tǒng)一時(shí)區(qū)為 UTC,避免時(shí)間處理歧義
def build_spark():
if os.path.exists(WAREHOUSE):
shutil.rmtree(WAREHOUSE)
return SparkSession.builder \
.appName("IcebergWorkshop") \
.config("spark.jars.packages", "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.8.1") \
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config(f"spark.sql.catalog.{CATALOG}", "org.apache.iceberg.spark.SparkCatalog") \
.config(f"spark.sql.catalog.{CATALOG}.type", "hadoop") \
.config(f"spark.sql.catalog.{CATALOG}.warehouse", WAREHOUSE) \
.config("spark.sql.session.timeZone", "UTC") \
.getOrCreate()
建表與數(shù)據(jù)寫(xiě)入
Iceberg 表創(chuàng)建時(shí)需要指定表名、Schema 和分區(qū)策略。與傳統(tǒng) Hive 表不同,Iceberg 采用隱藏分區(qū)(Hidden Partitioning)機(jī)制:用戶(hù)寫(xiě)入原始時(shí)間戳字段,Iceberg 自動(dòng)根據(jù)分區(qū)規(guī)范將數(shù)據(jù)寫(xiě)入對(duì)應(yīng)的分區(qū)目錄,查詢(xún)時(shí)自動(dòng)進(jìn)行分區(qū)裁剪。
創(chuàng)建無(wú)分區(qū)表 customers
- 無(wú)分區(qū)表將所有數(shù)據(jù)存儲(chǔ)在單一目錄下,適合數(shù)據(jù)量較小或查詢(xún)不帶過(guò)濾條件的場(chǎng)景。
CREATE TABLE local.iceberg_db.customers (
customer_id INT,
name STRING,
email STRING,
country STRING,
registration_date DATE,
tier STRING
)
USING iceberg
[OK] customers table created (unpartitioned)創(chuàng)建隱藏分區(qū)表 orders
PARTITIONED BY (months(order_time))定義了按月分區(qū)。用 order_time 的月份自動(dòng)分區(qū),但是表里不會(huì)多出一列 “月份” 字段- Iceberg 的隱藏分區(qū)將
order_time轉(zhuǎn)換為order_time_month分區(qū)列,數(shù)據(jù)寫(xiě)入order_time_month=2024-01/等目錄。 - 查詢(xún)時(shí)無(wú)需關(guān)心分區(qū)目錄結(jié)構(gòu),Iceberg 自動(dòng)根據(jù)時(shí)間條件過(guò)濾分區(qū)。
CREATE TABLE local.iceberg_db.orders (
order_id INT,
customer_id INT,
product_name STRING,
category STRING,
quantity INT,
price DOUBLE,
order_time TIMESTAMP,
status STRING
)
USING iceberg
PARTITIONED BY (months(order_time))
[OK] orders table created (hidden partition: months(order_time))向 customers 插入 20 行數(shù)據(jù)
sql> SELECT * FROM local.iceberg_db.customers ORDER BY customer_id +-----------+-------------+--------------------+-------+-----------------+-------+ |customer_id|name |email |country|registration_date|tier | +-----------+-------------+--------------------+-------+-----------------+-------+ |1 |張偉 |zhangwei@example.com|CN |2023-01-15 |vip | |2 |李娜 |lina@example.com |CN |2023-02-20 |premium| |3 |John Smith |jsmith@example.com |US |2023-03-10 |free | |4 |Emily Davis |edavis@example.com |US |2023-04-05 |premium| |5 |James Wilson |jwilson@example.com |UK |2023-05-18 |free | |6 |田中太郎 |tanaka@example.com |JP |2023-06-01 |vip | |7 |王芳 |wangfang@example.com|CN |2023-06-22 |free | |8 |Sarah Brown |sbrown@example.com |UK |2023-07-14 |premium| |9 |Mike Johnson |mjohnson@example.com|US |2023-08-30 |free | |10 |佐藤花子 |sato@example.com |JP |2023-09-12 |premium| |11 |趙磊 |zhaolei@example.com |CN |2023-10-01 |vip | |12 |Alice Lee |alee@example.com |US |2023-11-20 |free | |13 |David Clark |dclark@example.com |UK |2024-01-05 |premium| |14 |陳靜 |chenjing@example.com|CN |2024-01-15 |free | |15 |Robert Taylor|rtaylor@example.com |US |2024-02-28 |vip | |16 |鈴木一郎 |suzuki@example.com |JP |2024-03-10 |free | |17 |黃麗 |huangli@example.com |CN |2024-03-22 |premium| |18 |Emma Thomas |ethomas@example.com |UK |2024-04-01 |free | |19 |林強(qiáng) |linqiang@example.com|CN |2024-04-15 |vip | |20 |Tom Harris |tharris@example.com |US |2024-05-01 |premium| +-----------+-------------+--------------------+-------+-----------------+-------+
向 orders 插入 50 行數(shù)據(jù)
orders_df = generate_orders(spark)
orders_df.writeTo(f"{DB}.orders").append()
count(local.iceberg_db.orders WHERE 1=1) = 50
>>> orders sample (first 10)
sql> SELECT * FROM local.iceberg_db.orders ORDER BY order_id LIMIT 10
+--------+-----------+-------------------+-----------+--------+------+-------------------+---------+
|order_id|customer_id|product_name |category |quantity|price |order_time |status |
+--------+-----------+-------------------+-----------+--------+------+-------------------+---------+
|1001 |1 |MacBook Pro 16 |Electronics|1 |2499.0|2024-01-05 09:23:00|completed|
|1002 |1 |AirPods Pro |Electronics|2 |249.0 |2024-01-05 09:25:00|completed|
|1003 |2 |Python編程入門(mén) |Books |1 |59.9 |2024-01-12 14:30:00|completed|
|1004 |3 |數(shù)據(jù)密集型應(yīng)用設(shè)計(jì) |Books |1 |79.0 |2024-01-18 11:00:00|completed|
|1005 |4 |Sony WH-1000XM5 |Electronics|1 |349.0 |2024-01-25 16:45:00|cancelled|
|1006 |5 |冬季羽絨服 |Clothing |1 |199.0 |2024-02-01 10:15:00|completed|
|1007 |6 |Nintendo Switch |Electronics|1 |299.0 |2024-02-08 08:30:00|completed|
|1008 |7 |Spark快速大數(shù)據(jù)分析|Books |1 |69.0 |2024-02-14 13:20:00|pending |
|1009 |8 |Yoga Mat Premium |Home |2 |45.5 |2024-02-20 15:00:00|completed|
|1010 |9 |機(jī)械鍵盤(pán) Cherry軸 |Electronics|1 |129.0 |2024-02-28 09:45:00|cancelled|
+--------+-----------+-------------------+-----------+--------+------+-------------------+---------+orders 表按 months(order_time) 分區(qū),數(shù)據(jù)分布在 2024 年 1-6 月的 6 個(gè)分區(qū)目錄中。Parquet 文件存儲(chǔ)在 /tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-01/ 等目錄下。
查詢(xún)與分區(qū)裁剪
Iceberg 的分區(qū)裁剪在查詢(xún)規(guī)劃階段完成,通過(guò)讀取元數(shù)據(jù)文件中的分區(qū)統(tǒng)計(jì)信息,過(guò)濾掉不相關(guān)的分區(qū)目錄。與 Hive 的物理分區(qū)裁剪不同,Iceberg 無(wú)需掃描分區(qū)目錄,metadata-driven 方式效率更高。
分區(qū)裁剪驗(yàn)證
- 查詢(xún)條件
order_time >= '2024-06-01'觸發(fā)分區(qū)裁剪,Iceberg 只讀取order_time_month=2024-06分區(qū)的數(shù)據(jù)文件,跳過(guò) 1-5 月的分區(qū)。Hive 表需要依賴(lài)物理目錄結(jié)構(gòu)(order_time_month=2024-06/),Iceberg 的元數(shù)據(jù)驅(qū)動(dòng)方式更加靈活。
SELECT order_id, product_name, order_time
FROM local.iceberg_db.orders
WHERE order_time >= TIMESTAMP '2024-06-01 00:00:00'
ORDER BY order_time
>>> June 2024 orders (partition pruning)
sql>
SELECT order_id, product_name, order_time
FROM local.iceberg_db.orders
WHERE ORDER_TIME >= TIMESTAMP '2024-06-01 00:00:00'
ORDER BY order_time
+--------+----------------------+-------------------+
|order_id|product_name |order_time |
+--------+----------------------+-------------------+
|1038 |USB-C Hub 7合1 |2024-06-01 08:15:00|
|1039 |T恤 夏季純棉 |2024-06-03 10:00:00|
|1040 |Effective Java |2024-06-05 15:30:00|
|1041 |運(yùn)動(dòng)短褲 |2024-06-08 09:00:00|
|1042 |Air Conditioner |2024-06-10 11:00:00|
|1043 |Water Bottle Insulated|2024-06-12 13:00:00|
|1044 |算法導(dǎo)論 |2024-06-15 14:30:00|
|1045 |Wireless Mouse |2024-06-17 10:00:00|
|1046 |Bluetooth Headset |2024-06-19 16:00:00|
|1047 |Phone Case Premium |2024-06-21 08:30:00|
|1048 |Scala編程 |2024-06-23 11:15:00|
|1049 |Laptop Stand |2024-06-25 13:45:00|
|1050 |Pillow Memory Foam |2024-06-28 15:00:00|
+--------+----------------------+-------------------+UPDATE 與 DELETE
Iceberg 支持行級(jí) UPDATE 和 DELETE 操作,這是傳統(tǒng) Parquet/ORC 文件無(wú)法做到的能力。Iceberg 通過(guò) Copy-on-Write(寫(xiě)入時(shí)復(fù)制)或 Merge-on-Read(讀取時(shí)合并)策略實(shí)現(xiàn)行級(jí)操作。執(zhí)行 UPDATE/DELETE 后會(huì)創(chuàng)建新快照,舊數(shù)據(jù)文件保留用于時(shí)間旅行。
UPDATE 修改客戶(hù)名
- UPDATE 操作創(chuàng)建新快照,舊數(shù)據(jù)保留在歷史快照中??梢酝ㄟ^(guò)時(shí)間旅行查詢(xún)更新前的數(shù)據(jù)。
UPDATE local.iceberg_db.customers SET name = '張偉(已更名)' WHERE customer_id = 1 >>> BEFORE update sql> SELECT customer_id, name FROM local.iceberg_db.customers WHERE customer_id = 1 +-----------+----+ |customer_id|name| +-----------+----+ |1 |張偉| +-----------+----+ >>> AFTER update sql> SELECT customer_id, name FROM local.iceberg_db.customers WHERE customer_id = 1 +-----------+------------+ |customer_id|name | +-----------+------------+ |1 |張偉(已更名)| +-----------+------------+
DELETE 刪除已取消訂單
- 刪除了 5 條狀態(tài)為 cancelled 的訂單。DELETE 操作同樣創(chuàng)建新快照,舊數(shù)據(jù)可追溯。
- Iceberg 的行級(jí)操作能力使其適合 GDPR 合規(guī)刪除和實(shí)時(shí)數(shù)據(jù)更新場(chǎng)景。
DELETE FROM local.iceberg_db.orders WHERE status = 'cancelled'
>>> cancelled orders (BEFORE delete)
sql> SELECT order_id, status FROM local.iceberg_db.orders WHERE status = 'cancelled'
+--------+---------+
|order_id|status |
+--------+---------+
|1005 |cancelled|
|1010 |cancelled|
|1032 |cancelled|
|1045 |cancelled|
|1018 |cancelled|
+--------+---------+
count(local.iceberg_db.orders WHERE status = 'cancelled') = 0
[OK] All cancelled orders deleted, verified count = 0Schema 演進(jìn)
Iceberg 的 Schema 演進(jìn)是 metadata-only 操作,不涉及數(shù)據(jù)文件重寫(xiě)。Iceberg 為每個(gè)列分配唯一的 column ID,Schema 演進(jìn)只更新元數(shù)據(jù)中的列映射關(guān)系,歷史數(shù)據(jù)文件的列 ID 保持不變。這使得 Iceberg 可以安全地進(jìn)行 Schema 演進(jìn)而不破壞歷史數(shù)據(jù)。
ADD COLUMN 添加列
ALTER TABLE local.iceberg_db.customers ADD COLUMN phone STRING ALTER TABLE local.iceberg_db.customers ADD COLUMN loyalty_points INT DEFAULT 0 >>> Schema after ADD COLUMN sql> DESCRIBE local.iceberg_db.customers +-----------------+---------+-------+ |col_name |data_type|comment| +-----------------+---------+-------+ |customer_id |int |NULL | |name |string |NULL | |email |string |NULL | |country |string |NULL | |registration_date|date |NULL | |tier |string |NULL | |phone |string |NULL | |loyalty_points |int |NULL | +-----------------+---------+-------+
新增 phone 和 loyalty_points 兩列。DEFAULT 0 表示新列的默認(rèn)值,寫(xiě)入時(shí)不指定該字段會(huì)自動(dòng)填充。
RENAME COLUMN 重命名列
ALTER TABLE local.iceberg_db.customers RENAME COLUMN name TO full_name >>> Schema after RENAME name → full_name sql> DESCRIBE local.iceberg_db.customers +-----------------+---------+-------+ |col_name |data_type|comment| +-----------------+---------+-------+ |customer_id |int |NULL | |full_name |string |NULL | |email |string |NULL | |country |string |NULL | |registration_date|date |NULL | |tier |string |NULL | |phone |string |NULL | |loyalty_points |int |NULL | +-----------------+---------+-------+
列名從 name 變更為 full_name,數(shù)據(jù)文件中的內(nèi)容不受影響。Iceberg 內(nèi)部通過(guò) column ID 追蹤列,rename 只更新元數(shù)據(jù)映射。
ALTER TYPE 類(lèi)型提升
ALTER TABLE local.iceberg_db.customers ALTER COLUMN loyalty_points TYPE BIGINT >>> Schema after INT → BIGINT promotion sql> DESCRIBE local.iceberg_db.customers +-----------------+---------+-------+ |col_name |data_type|comment| +-----------------+---------+-------+ |customer_id |int |NULL | |full_name |string |NULL | |email |string |NULL | |country |string |NULL | |registration_date|date |NULL | |tier |string |NULL | |phone |string |NULL | |loyalty_points |bigint |NULL | +-----------------+---------+-------+
loyalty_points 從 INT 提升為 BIGINT。Iceberg 支持安全的類(lèi)型提升(int→bigint,float→double),不兼容的類(lèi)型變更會(huì)被拒絕。
DROP COLUMN 刪除列
ALTER TABLE local.iceberg_db.customers DROP COLUMN registration_date >>> Schema after DROP registration_date sql> DESCRIBE local.iceberg_db.customers +--------------+---------+-------+ |col_name |data_type|comment| +--------------+---------+-------+ |customer_id |int |NULL | |full_name |string |NULL | |email |string |NULL | |country |string |NULL | |tier |string |NULL | |phone |string |NULL | |loyalty_points|bigint |NULL | +--------------+---------+-------+
刪除了 registration_date 列。DROP COLUMN 是 metadata-only 操作,歷史數(shù)據(jù)文件中該列的數(shù)據(jù)仍然存在,只是查詢(xún)時(shí)不再返回。
REORDER COLUMN 調(diào)整列順序
ALTER TABLE local.iceberg_db.customers ALTER COLUMN email FIRST >>> Schema after MOVE email FIRST sql> DESCRIBE local.iceberg_db.customers +--------------+---------+-------+ |col_name |data_type|comment| +--------------+---------+-------+ |email |string |NULL | |customer_id |int |NULL | |full_name |string |NULL | |country |string |NULL | |tier |string |NULL | |phone |string |NULL | |loyalty_points|bigint |NULL | +--------------+---------+-------+
email 列被移動(dòng)到表的第一位。列順序調(diào)整不影響數(shù)據(jù)存儲(chǔ),只是改變查詢(xún)結(jié)果的顯示順序。
Schema 演進(jìn)后數(shù)據(jù)驗(yàn)證
SELECT customer_id, full_name, email, phone, loyalty_points FROM local.iceberg_db.customers ORDER BY customer_id LIMIT 5 >>> Data still intact after schema evolution sql> SELECT customer_id, full_name, email, phone, loyalty_points FROM local.iceberg_db.customers ORDER BY customer_id LIMIT 5 +-----------+------------+--------------------+-----+--------------+ |customer_id|full_name |email |phone|loyalty_points| +-----------+------------+--------------------+-----+--------------+ |1 |張偉(已更名)|zhangwei@example.com|NULL |NULL | |2 |李娜 |lina@example.com |NULL |NULL | |3 |John Smith |jsmith@example.com |NULL |NULL | |4 |Emily Davis |edavis@example.com |NULL |NULL | |5 |James Wilson|jwilson@example.com |NULL |NULL | +-----------+------------+--------------------+-----+--------------+
經(jīng)過(guò)前述的所有 Schema 演進(jìn)操作后,原有數(shù)據(jù)完整保留。新增的 phone 和 loyalty_points 列為 NULL(未賦值),UPDATE 后的 full_name 也正確保留。這驗(yàn)證了 Iceberg Schema 演進(jìn)的數(shù)據(jù)兼容性。
元數(shù)據(jù)表探查
Iceberg 表的元數(shù)據(jù)層包含四級(jí)結(jié)構(gòu):metadata file → manifest list → manifest → data file。Iceberg 提供元數(shù)據(jù)表(metadata tables)供用戶(hù)直接查詢(xún)這些元數(shù)據(jù),無(wú)需訪問(wèn)底層文件。
$snapshots:列出所有快照及操作摘要$history:快照時(shí)間線(xiàn)及父子關(guān)系$files:列出所有數(shù)據(jù)文件$manifests:列出所有 manifest 文件
$snapshots 顯示兩個(gè)快照:第一個(gè)是初始 append(50 條記錄),第二個(gè)是 DELETE 后的 overwrite(刪除 5 條,剩余 45 條)
- spark.app.id:local-1779244843156,Spark 應(yīng)用 ID,本地運(yùn)行的 Spark 任務(wù)
- added-data-files:6,本次新增 6 個(gè)數(shù)據(jù)文件
- added-records:50,本次寫(xiě)入 50 條數(shù)據(jù)
- added-files-size:17455,新增文件總大小 17455 字節(jié)(約 17KB)
- changed-partition-count:6,本次操作影響 6 個(gè)分區(qū)
- total-records:50,表總數(shù)據(jù)量 50 條(首次寫(xiě)入)
- total-files-size:17455,表總大小 17455 字節(jié)(約 17KB)
- total-data-files:6,表總數(shù)據(jù)文件 6 個(gè)
- total-delete-files:0,表總刪除文件 0 個(gè)
- total-position-deletes:0,位置刪除數(shù)據(jù) 0 條
- total-equality-deletes:0,等值刪除數(shù)據(jù) 0 條
- engine-name:spark,計(jì)算引擎為 Spark
- engine-version:3.5.6,Spark 版本 3.5.6
- iceberg-version:Apache Iceberg 1.8.1,Iceberg 版本 1.8.1
SELECT snapshot_id, committed_at, operation, summary
FROM local.iceberg_db.orders.snapshots
ORDER BY committed_at
>>> snapshots
|snapshot_id |committed_at |operation|summary
|2262767112950789139|2026-05-20 02:40:54.961|append|{spark.app.id -> local-1779244843156, added-data-files -> 6, added-records -> 50, added-files-size -> 17455, changed-partition-count -> 6, total-records -> 50, total-files-size -> 17455, total-data-files -> 6, total-delete-files -> 0, total-position-deletes -> 0, total-equality-deletes -> 0, engine-version -> 3.5.6, app-id -> local-1779244843156, engine-name -> spark, iceberg-version -> Apache Iceberg 1.8.1 (commit 9ce0fcf0af7becf25ad9fc996c3bad2afdcfd33d)}|
|5067017495298132779|2026-05-20 02:40:58.261|overwrite|{spark.app.id -> local-1779244843156, added-data-files -> 5, deleted-data-files -> 5, added-records -> 38, deleted-records -> 43, added-files-size -> 14369, removed-files-size -> 14642, changed-partition-count -> 5, total-records -> 45, total-files-size -> 17182, total-data-files -> 6, total-delete-files -> 0, total-position-deletes -> 0, total-equality-deletes -> 0, engine-version -> 3.5.6, app-id -> local-1779244843156, engine-name -> spark, iceberg-version -> Apache Iceberg 1.8.1 (commit 9ce0fcf0af7becf25ad9fc9963bad2afdcfd33d)}|$history 展示快照的父子鏈,第一個(gè)快照無(wú) parent(根快照)
SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor FROM local.iceberg_db.orders.history ORDER BY made_current_at >>> history sql> SELECT made_current_at, snapshot_id, parent_id, is_current_ancestor FROM local.iceberg_db.orders.history ORDER BY made_current_at +-----------------------+-------------------+-------------------+-------------------+ |made_current_at |snapshot_id |parent_id |is_current_ancestor| +-----------------------+-------------------+-------------------+-------------------+ |2026-05-20 02:40:54.961|2262767112950789139|NULL |true | |2026-05-20 02:40:58.261|5067017495298132779|2262767112950789139|true | +-----------------------+-------------------+-------------------+-------------------+
$files 列出 6 個(gè) Parquet 文件,按月份分區(qū)分布
SELECT content, file_path, file_format, record_count FROM local.iceberg_db.orders.files >>> data files sql> SELECT content, file_path, file_format, record_count FROM local.iceberg_db.orders.files |content|file_path |file_format|record_count| |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-01/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00001.parquet|PARQUET |4 | |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-02/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00002.parquet|PARQUET |4 | |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-05/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00004.parquet|PARQUET |10 | |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-06/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00003.parquet|PARQUET |12 | |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-04/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00003.parquet|PARQUET |8 | |0 |/tmp/iceberg-warehouse/iceberg_db/orders/data/order_time_month=2024-03/00000-30-1f1b8605-71f6-4e42-bff3-59720a357572-0-00001.parquet|PARQUET |7 |
$manifests 列出 2 個(gè) manifest 文件,記錄了數(shù)據(jù)文件的元信息
SELECT path, length, partition_spec_id FROM local.iceberg_db.orders.manifests >>> manifests sql> SELECT path, length, partition_spec_id FROM local.iceberg_db.orders.manifests |path |length|partition_spec_id| |/tmp/iceberg-warehouse/iceberg_db/orders/metadata/46ac57f8-4ee9-4070-9a0e-7b7c56e33ee6-m1.avro|8278|0| |/tmp/iceberg-warehouse/iceberg_db/orders/metadata/46ac57f8-4ee9-4070-9a0e-7b7c56e33ee6-m0.avro|8420|0|
時(shí)間旅行與快照回滾
Iceberg 的 MVCC 快照模型支持時(shí)間旅行查詢(xún)。每個(gè)快照有唯一的 snapshot_id 和 commit 時(shí)間戳,可以隨時(shí)回溯到歷史狀態(tài)。VERSION AS OF 通過(guò)快照 ID 查詢(xún),TIMESTAMP AS OF 通過(guò)時(shí)間戳查詢(xún)?;貪L操作是 O(1) 的元數(shù)據(jù)修改,不涉及數(shù)據(jù)復(fù)制。
VERSION AS OF 按快照 ID 查詢(xún)
- 第一個(gè)快照包含 50 條訂單(DELETE 之前的數(shù)據(jù))。通過(guò)快照 ID 可以精確訪問(wèn)歷史狀態(tài)。
SELECT count(*) as cnt FROM local.iceberg_db.orders VERSION AS OF 2262767112950789139 >>> Count at first snapshot (before DELETE) sql> SELECT count(*) as cnt FROM local.iceberg_db.orders VERSION AS OF 2262767112950789139 +---+ |cnt| +---+ |50 | +---+
TIMESTAMP AS OF 按時(shí)間戳查詢(xún)
- 時(shí)間戳查詢(xún)返回該時(shí)刻對(duì)應(yīng)的快照數(shù)據(jù)。Iceberg 自動(dòng)找到最接近該時(shí)間戳的快照進(jìn)行查詢(xún)。
SELECT count(*) as cnt FROM local.iceberg_db.orders FOR TIMESTAMP AS OF '2026-05-20 02:40:54.961000' >>> Count at timestamp 2026-05-20 02:40:54.961000 sql> SELECT count(*) as cnt FROM local.iceberg_db.orders FOR TIMESTAMP AS OF '2026-05-20 02:40:54.961000' +---+ |cnt| +---+ |50 | +---+
rollback_to_snapshot 快照回滾
CALL local.system.rollback_to_snapshot('iceberg_db.orders', 2262767112950789139)
Rolling back to snapshot: 2262767112950789139
>>> Count after rollback
sql> SELECT count(*) as cnt FROM local.iceberg_db.orders
+---+
|cnt|
+---+
|50 |
+---+回滾到第一個(gè)快照后,表狀態(tài)恢復(fù)到 DELETE 之前?;貪L是元數(shù)據(jù)操作,速度極快,不涉及數(shù)據(jù)文件復(fù)制。
從舊快照恢復(fù)被刪除的數(shù)據(jù)
- 從歷史快照查詢(xún)被刪除的 cancelled 訂單,重新插入當(dāng)前表。注意這里恢復(fù)出 10 條(包含之前測(cè)試過(guò)程中累積的),因?yàn)榛貪L后再次執(zhí)行 MERGE 等操作增加了數(shù)據(jù)。
INSERT INTO local.iceberg_db.orders SELECT * FROM local.iceberg_db.orders VERSION AS OF 2262767112950789139 WHERE status = 'cancelled' >>> Cancelled orders restored sql> SELECT count(*) as cnt FROM local.iceberg_db.orders WHERE status = 'cancelled' +---+ |cnt| +---+ |10 | +---+
驗(yàn)證回滾和恢復(fù)后的表狀態(tài)
SELECT status, count(*) as cnt
FROM local.iceberg_db.orders
GROUP BY status
ORDER BY status
count(local.iceberg_db.orders WHERE 1=1) = 55
>>> Order status distribution
sql> SELECT status, count(*) as cnt FROM local.iceberg_db.orders GROUP BY status ORDER BY status
+---------+---+
|status |cnt|
+---------+---+
|cancelled|10 |
|completed|41 |
|pending |4 |
+---------+---+表現(xiàn)在有 55 條記錄,包括 10 條已恢復(fù)的 cancelled 訂單。時(shí)間旅行和回滾能力使得數(shù)據(jù)恢復(fù)變得簡(jiǎn)單可靠。
分區(qū)演進(jìn)
分區(qū)演進(jìn)允許修改已有表的分區(qū)策略,無(wú)需重寫(xiě)歷史數(shù)據(jù)。隨著數(shù)據(jù)量增長(zhǎng)或查詢(xún)模式變化,可能需要從細(xì)粒度分區(qū)(如按月)調(diào)整為粗粒度分區(qū)(如按年)。Iceberg 的分區(qū)規(guī)范與數(shù)據(jù)文件分離存儲(chǔ),舊數(shù)據(jù)保持原有分區(qū)規(guī)范,新數(shù)據(jù)使用新規(guī)范,查詢(xún)時(shí)自動(dòng)適配。
ALTER TABLE REPLACE PARTITION FIELD
ALTER TABLE local.iceberg_db.orders REPLACE PARTITION FIELD months(order_time) WITH years(order_time) >>> Partition spec after evolution (look for Partition Spec) sql> DESCRIBE EXTENDED local.iceberg_db.orders +----------------------------+------------------------------------------------+-------+ |col_name |data_type |comment| +----------------------------+------------------------------------------------+-------+ |order_id |int |NULL | |customer_id |int |NULL | |product_name |string |NULL | |category |string |NULL | |quantity |int |NULL | |price |double |NULL | |order_time |timestamp |NULL | |status |string |NULL | | | | | |# Partitioning | | | |Part 0 |years(order_time) | | | | | | |# Metadata Columns | | | |_spec_id |int |NULL | |_partition |struct<order_time_month:int,order_time_year:int>|NULL | |_file |string |NULL | |_pos |bigint |NULL | |_deleted |boolean |NULL | | | | | |# Detailed Table Information| | | +----------------------------+------------------------------------------------+-------+ only showing top 20 rows
分區(qū)規(guī)范從 months(order_time) 變?yōu)?years(order_time)。注意 _partition 元數(shù)據(jù)列同時(shí)包含 order_time_month 和 order_time_year,說(shuō)明 Iceberg 保留了歷史分區(qū)信息以支持對(duì)舊數(shù)據(jù)的透明查詢(xún)。Hive 表的分區(qū)變更通常需要數(shù)據(jù)遷移,Iceberg 的分區(qū)演進(jìn)更加靈活。
MERGE INTO 與 CDC
MERGE INTO 是 Iceberg 支持的原子 upsert 操作,類(lèi)似于 “INSERT … ON CONFLICT UPDATE”。這對(duì)于 CDC(Change Data Capture)場(chǎng)景非常有用:定期從上游系統(tǒng)接收變更數(shù)據(jù)流,通過(guò) MERGE INTO 同步到 Iceberg 表,支持增量更新和插入。
MERGE INTO 實(shí)現(xiàn) CDC 同步
首先創(chuàng)建 staging 表作為 CDC 數(shù)據(jù)源:
CREATE TABLE local.iceberg_db.orders_staging (
order_id INT, customer_id INT, product_name STRING,
category STRING, quantity INT, price DOUBLE,
order_time TIMESTAMP,
status STRING
) USING iceberg插入 CDC 數(shù)據(jù):
cdc_rows = [
(1001, 1, "MacBook Pro 16 (M4)", "Electronics", 1, 2799.00, dt(2024,1,5,9,23), "completed"),
(1051, 2, "Pixel Watch 3", "Electronics", 1, 349.00, dt(2024,6,30,10,0), "pending"),
(1052, 5, "Standing Desk", "Home", 1, 499.00, dt(2024,6,30,11,0), "completed"),
]
cdc_df = spark.createDataFrame([Row(*r) for r in cdc_rows], schema=...)
cdc_df.writeTo(f"{DB}.orders_staging").append()查看合并前的數(shù)據(jù):
>>> BEFORE MERGE: target rows sql> SELECT order_id, product_name, price, status FROM local.iceberg_db.orders WHERE order_id IN (1001, 1051, 1052) ORDER BY order_id +--------+--------------+------+---------+ |order_id|product_name |price |status | +--------+--------------+------+---------+ |1001 |MacBook Pro 16|2499.0|completed| +--------+--------------+------+---------+ >>> BEFORE MERGE: source (CDC) rows sql> SELECT order_id, product_name, price, status FROM local.iceberg_db.orders_staging ORDER BY order_id +--------+-------------------+------+---------+ |order_id|product_name |price |status | +--------+-------------------+------+---------+ |1001 |MacBook Pro 16 (M4)|2799.0|completed| |1051 |Pixel Watch 3 |349.0 |pending | |1052 |Standing Desk |499.0 |completed| +--------+-------------------+------+---------+
執(zhí)行 MERGE INTO:
MERGE INTO local.iceberg_db.orders t USING local.iceberg_db.orders_staging s ON t.order_id = s.order_id WHEN MATCHED THEN UPDATE SET * # 用源表的所有字段覆蓋目標(biāo)表 WHEN NOT MATCHED THEN INSERT * # 新數(shù)據(jù)直接插入目標(biāo)表
查看合并后的結(jié)果:
- order_id=1001 已存在,執(zhí)行 UPDATE(價(jià)格從 2499.0 更新為 2799.0)
- order_id=1051 和 1052 不存在,執(zhí)行 INSERT
- MERGE INTO 是原子操作,保證數(shù)據(jù)一致性
>>> AFTER MERGE: updated + inserted rows sql> SELECT order_id, product_name, price, status FROM local.iceberg_db.orders WHERE order_id IN (1001, 1051, 1052) ORDER BY order_id +--------+-------------------+------+---------+ |order_id|product_name |price |status | +--------+-------------------+------+---------+ |1001 |MacBook Pro 16 (M4)|2799.0|completed| |1051 |Pixel Watch 3 |349.0 |pending | |1052 |Standing Desk |499.0 |completed| +--------+-------------------+------+---------+
這是典型的 CDC 同步模式,staging 表存儲(chǔ)來(lái)自上游的變更數(shù)據(jù),通過(guò) MERGE INTO 同步到目標(biāo)表。
表維護(hù)
Iceberg 表在長(zhǎng)期運(yùn)行中會(huì)產(chǎn)生三類(lèi)維護(hù)問(wèn)題:
- 小文件問(wèn)題:頻繁的小規(guī)模寫(xiě)入導(dǎo)致大量小文件,影響查詢(xún)性能
- 快照膨脹:每個(gè)寫(xiě)入操作產(chǎn)生新快照,歷史快照占用空間
- 孤兒文件:compaction 或其他操作失敗后遺留的未引用文件
Iceberg 提供系統(tǒng)存儲(chǔ)過(guò)程進(jìn)行維護(hù):
rewrite_data_files(壓縮)expire_snapshots(過(guò)期快照)remove_orphan_files(清理孤兒文件)
制造碎片文件
for i in range(10):
spark.sql(f"""
INSERT INTO {DB}.orders
SELECT 1060 + {i}, 1, 'Test Product {i}', 'Home', 1, {9.90 + i},
TIMESTAMP '2024-07-0{(i%9)+1} 10:00:00', 'completed'
""")
Data files BEFORE compaction: 21通過(guò)循環(huán)插入 10 條記錄,創(chuàng)建了 10 個(gè)新的小文件。加上原有的 11 個(gè)文件,總計(jì) 21 個(gè)數(shù)據(jù)文件。
Compaction 壓縮數(shù)據(jù)文件。rewrite_data_files 將 21 個(gè)小文件合并為 2 個(gè)大文件(目標(biāo)大小 128MB)。Compaction 顯著提升查詢(xún)性能,減少文件元數(shù)據(jù)開(kāi)銷(xiāo)。
CALL local.system.rewrite_data_files(
table => 'iceberg_db.orders',
options => map('target-file-size-bytes', '134217728')
)
Data files AFTER compaction: 2
Files reduced: 21 → 2expire_snapshots 過(guò)期快照。retain_last=3 保留最近 3 個(gè)快照,其余 12 個(gè)快照被過(guò)期。older_than 設(shè)置為遙遠(yuǎn)的未來(lái)時(shí)間,確保只按 retain_last 參數(shù)過(guò)期。過(guò)期快照后,其關(guān)聯(lián)的數(shù)據(jù)文件如果不再被其他快照引用,將成為孤兒文件。
CALL local.system.expire_snapshots(
table => 'iceberg_db.orders',
older_than => TIMESTAMP '2099-01-01 00:00:00',
retain_last => 3
)
Snapshots BEFORE expiration: 15
Snapshots AFTER expiration (retain_last=3): 3remove_orphan_files 清理孤兒文件。清理因 compaction 和快照過(guò)期產(chǎn)生的孤兒文件。建議定期執(zhí)行維護(hù)任務(wù)(如每天或每周),保持表健康。
CALL local.system.remove_orphan_files(
table => 'iceberg_db.orders'
)
[OK] Orphan file cleanup done視圖與總結(jié)
創(chuàng)建臨時(shí)視圖
- Hadoop Catalog 不支持持久化視圖,使用臨時(shí)視圖代替。視圖是對(duì) Iceberg 表的查詢(xún)抽象,可以簡(jiǎn)化復(fù)雜分析。
CREATE TEMP VIEW v_order_summary AS
SELECT
year(order_time) AS order_year,
month(order_time) AS order_month,
category,
count(*) AS order_count,
sum(quantity) AS total_quantity,
round(sum(price * quantity), 2) AS total_revenue
FROM local.iceberg_db.orders
GROUP BY year(order_time), month(order_time), category
ORDER BY order_year, order_month, category
NOTE: Hadoop Catalog does not support persistent views, using temp view instead
>>> Temp view query result
sql> SELECT * FROM v_order_summary
+----------+-----------+-----------+-----------+--------------+-------------+
|order_year|order_month|category |order_count|total_quantity|total_revenue|
+----------+-----------+-----------+-----------+--------------+-------------+
|2024 |1 |Books |2 |2 |138.9 |
|2024 |1 |Electronics|4 |5 |3995.0 |
|2024 |2 |Books |1 |1 |69.0 |
|2024 |2 |Clothing |1 |1 |199.0 |
|2024 |2 |Electronics|3 |3 |557.0 |
|2024 |2 |Home |1 |2 |91.0 |
|2024 |3 |Books |2 |2 |164.0 |
|2024 |3 |Clothing |1 |1 |129.0 |
|2024 |3 |Electronics|3 |3 |1547.0 |
|2024 |3 |Home |1 |1 |89.99 |
|2024 |4 |Books |4 |4 |326.0 |
|2024 |4 |Clothing |1 |2 |316.0 |
|2024 |4 |Electronics|4 |4 |1056.0 |
|2024 |4 |Home |1 |1 |35.0 |
|2024 |5 |Books |4 |4 |258.0 |
|2024 |5 |Clothing |1 |3 |179.7 |
|2024 |5 |Electronics|4 |4 |1546.0 |
|2024 |5 |Home |3 |4 |538.0 |
|2024 |6 |Books |3 |3 |269.0 |
|2024 |6 |Clothing |2 |3 |128.8 |
+----------+-----------+-----------+-----------+--------------+-------------+
only showing top 20 rows到此這篇關(guān)于基于Pyspark對(duì)Apache Iceberg核心功能的使用實(shí)踐指南的文章就介紹到這了,更多相關(guān)Pyspark Apache Iceberg使用內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Python利用PySpark和Kafka實(shí)現(xiàn)流處理引擎構(gòu)建指南
- Python大數(shù)據(jù)分析之PySpark原理與實(shí)戰(zhàn)教程詳解
- 使用用Pyspark和GraphX實(shí)現(xiàn)解析復(fù)雜網(wǎng)絡(luò)數(shù)據(jù)
- 從Pyspark UDF調(diào)用另一個(gè)自定義Python函數(shù)的方法步驟
- Python?PySpark案例實(shí)戰(zhàn)教程
- pandas與pyspark計(jì)算效率對(duì)比分析
- pyspark?dataframe列的合并與拆分實(shí)例
- PySpark中RDD的數(shù)據(jù)輸出問(wèn)題詳解
- pyspark自定義UDAF函數(shù)調(diào)用報(bào)錯(cuò)問(wèn)題解決
相關(guān)文章
Linux系統(tǒng)使用用戶(hù)密鑰ssh主機(jī)訪問(wèn)
這篇文章主要介紹了Linux系統(tǒng)使用用戶(hù)密鑰ssh主機(jī)訪問(wèn),它在安全上完全大于直接輸入root 的密碼,有需要的可以了解一下。2016-10-10
Linux系統(tǒng)安裝NoSQL(MongoDB和Redis)步驟及問(wèn)題解決辦法(總結(jié)篇)
這篇文章主要介紹了Linux系統(tǒng)安裝NoSQL(MongoDB和Redis)步驟及問(wèn)題解決辦法的相關(guān)資料,本文分步驟給大家介紹的非常詳細(xì),具有參考借鑒價(jià)值,感興趣的朋友一起看看吧2016-10-10
linux虛擬網(wǎng)絡(luò)設(shè)備之vlan配置詳解
這篇文章主要給大家介紹了關(guān)于linux虛擬網(wǎng)絡(luò)設(shè)備之vlan配置的相關(guān)資料,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧。2017-12-12
用vnc實(shí)現(xiàn)Windows遠(yuǎn)程連接linux桌面之服務(wù)器配置
這篇文章主要介紹了用vnc實(shí)現(xiàn)Windows遠(yuǎn)程連接linux桌面之服務(wù)器配置,需要的朋友可以參考下2016-09-09

