最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

基于Pyspark對(duì)Apache Iceberg核心功能的使用實(shí)踐指南

 更新時(shí)間:2026年05月21日 10:03:11   作者:zhaojiew10  
本文測(cè)試了Apache Iceberg表格式的核心功能,使用PySpark 3.5.6和Hadoop Catalog搭建本地實(shí)驗(yàn)環(huán)境,Iceberg作為開(kāi)放表格式,提供ACID事務(wù)、隱藏分區(qū)、Schema演進(jìn)等數(shù)據(jù)庫(kù)級(jí)特性,本文介紹基于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-warehouse
  • spark.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 = 0

Schema 演進(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   |
+-----------------+---------+-------+

新增 phoneloyalty_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ù)完整保留。新增的 phoneloyalty_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_monthorder_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)題:

  1. 小文件問(wèn)題:頻繁的小規(guī)模寫(xiě)入導(dǎo)致大量小文件,影響查詢(xún)性能
  2. 快照膨脹:每個(gè)寫(xiě)入操作產(chǎn)生新快照,歷史快照占用空間
  3. 孤兒文件: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 → 2

expire_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): 3

remove_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)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評(píng)論

临泽县| 芜湖市| 富源县| 教育| 湾仔区| 商河县| 巫溪县| 广元市| 兴安县| 柳河县| 叶城县| 剑河县| 奉节县| 泸水县| 井陉县| 澄迈县| 延庆县| 新郑市| 陆河县| 嵊泗县| 洪江市| 荃湾区| 汉寿县| 镇巴县| 磴口县| 石河子市| 崇明县| 册亨县| 烟台市| 光山县| 来凤县| 古交市| 榆社县| 天水市| 黎川县| 丹巴县| 曲阳县| 曲松县| 河西区| 中超| 察雅县|