AnalyticDB for MySQL深度整合Apache Iceberg等開源湖表格式,基於自研高效能XIHE引擎以及託管Spark引擎,提供開放、多引擎相容的資料湖(Lakehouse)能力。使用者建立的湖表資料以標準Parquet 格式持久化至Object Storage Service(阿里雲OSS),任何支援Iceberg的計算引擎(如 Spark、Flink、Trino)均可直接讀取,避免廠商鎖定,保障資料資產長期可用性。本文以Apache Iceberg和Delta Lake格式為例,介紹資料湖表的建立方法。
儲存託管模式
模式選擇
在建立資料湖表前,您需要決定資料的儲存位置。AnalyticDB for MySQL提供兩種靈活的儲存管理員模式,兼顧控制權與便捷性,滿足不同安全與營運需求:
使用者自有OSS Bucket
資料完全存放在使用者指定的同地區OSS Bucket中,滿足強合規與資料主權要求。建庫、建表時需顯式聲明儲存路徑,實現細粒度管控。
AnalyticDB for MySQL託管湖儲存
資料由AnalyticDB for MySQL自動管理底層儲存桶,該儲存桶在使用者帳號下不可見。使用者通過標準SQL無縫讀寫湖表,無需關心檔案系統、許可權配置或生命週期管理,極大降低使用門檻。詳情請參見湖儲存。
核心優勢
所有湖表的中繼資料(包括資料庫、表結構、列定義、分區資訊等)均由AnalyticDB for MySQL內建Catalog服務統一管理,無需使用者部署、擴縮容或維護獨立的中繼資料叢集。
使用開放資料湖如同操作傳統資料庫一樣簡單:
使用者只需關注
CREATE TABLE和SQL查詢邏輯。底層儲存、中繼資料、檔案格式、壓縮編碼等基礎設施由平台全託管。
保留對開源生態的完全開放性,實現 “簡單如資料庫,開放如資料湖” 的理想Lakehouse體驗。
前提條件
通過XIHE引擎建立Apache Iceberg表,叢集核心版本需為3.2.7及以上版本。
說明請在雲原生資料倉儲AnalyticDB MySQL控制台集群信息頁面,配寘資訊地區,查看和升級核心版本。若叢集已是最新預設基準版本但仍需升級,請通過DingTalk聯絡阿里雲服務支援處理(DingTalk帳號:
x5v_rm8wqzuqf)。在AnalyticDB for MySQL託管湖儲存中建立表,需提交工單聯絡支援人員開通湖儲存功能,並建立湖儲存。
建立預設表格式
您可以通過資料庫屬性設定預設的表格式,使該資料庫下建立的表自動採用指定格式,無需每次建表時單獨聲明。
在Spark Job資源群組、Interactive資源群組或提交的作業中,設定如下參數:
spark.sql.adb.sources.extractProviderFromDBProperties.enabled true建立資料庫時,通過
DBPROPERTIES指定'storage.format'為以下格式之一:delta、iceberg、parquet、orc。樣本如下:CREATE DATABASE IF NOT EXISTS db_storage_format LOCATION 'oss://path/to/db/' WITH DBPROPERTIES ('storage.format'='delta');執行以上語句後,在
db_storage_format庫下建立的表預設為delta類型。如您建表時通過using ${tableFormat}顯式指定表類型,則優先以顯式指定的表類型為準。
建立Apache Iceberg表
在使用者自有OSS bucket中建立表
建立非分區表
適用情境:小型維度資料表(如國家、地區、產品類別等)、全表掃描頻繁但資料量小、無需按時間或高基數欄位裁剪的待用資料。
XIHE SQL
CREATE DATABASE db_iceberg; -- 建立非分區的 nation 表(小維表,通常 < 100 行) CREATE TABLE db_iceberg.nation ( n_nationkey INT, n_name STRING, n_regionkey INT, n_comment STRING ) STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE DATABASE db_iceberg; CREATE TABLE db_iceberg.nation ( n_nationkey INT, n_name STRING, n_regionkey INT, n_comment STRING ) USING iceberg LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
建立分區表
AnalyticDB for MySQL支援使用分區轉換函式定義分區。分區表(Partitioned Tables) 是構建高效能、可擴充、易管理的現代資料湖(Lakehouse)的核心實踐之一。AnalyticDB for MySQL支援的Apache Iceberg分區轉換規則如下:
轉換函式 | 文法樣本 | 說明 | 適用類型 |
|
| 原值分區(等同於Hive分區) | 所有類型(但不推薦用於高基數欄位) |
|
| 按年分區 | Timestamp, Date |
|
| 按月分區 | Timestamp, Date |
|
| 按天分區(最常用) | Timestamp, Date |
|
| 按小時分區 | Timestamp |
|
| 雜湊分桶(N 為桶數) | 所有類型(常用於高基數欄位如 ID) |
|
| 截斷字串前 len 位 | String |
以XIHE、Spark引擎為例,不同情境適合的分區策略,樣本如下:
基於低基數類別欄位進行分區
XIHE
-- 按市場細分(mktsegment)原值分區,該欄位僅5個枚舉值,適合 identity CREATE TABLE db_iceberg.customer ( c_custkey BIGINT, c_name STRING, c_address STRING, c_nationkey INT, c_phone STRING, c_acctbal DECIMAL(15,2), c_mktsegment STRING, -- 低基數:'AUTOMOBILE', 'BUILDING', 'FURNITURE' 等 c_comment STRING ) PARTITIONED BY (c_mktsegment) STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.customer ( c_custkey BIGINT, c_name STRING, c_address STRING, c_nationkey INT, c_phone STRING, c_acctbal DECIMAL(15,2), c_mktsegment STRING, c_comment STRING ) USING iceberg PARTITIONED BY (c_mktsegment) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
按時間維度(years/months/days)分層分區
XIHE
-- 按訂單日期按天分區(最常用),兼顧查詢效能與管理成本 CREATE TABLE db_iceberg.orders ( o_orderkey BIGINT, o_custkey BIGINT, o_orderstatus STRING, o_totalprice DECIMAL(15,2), o_orderdate DATE, -- TPC-H 原生欄位 o_orderpriority STRING, o_clerk STRING, o_shippriority INT, o_comment STRING ) PARTITIONED BY (day(o_orderdate)) -- 推薦:平衡分區數量與裁剪效率 STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.orders ( o_orderkey BIGINT, o_custkey BIGINT, o_orderstatus STRING, o_totalprice DECIMAL(15,2), o_orderdate DATE, o_orderpriority STRING, o_clerk STRING, o_shippriority INT, o_comment STRING ) USING iceberg -- 注意:Iceberg 在 Spark 中使用的是 days() 轉換函式,且不需要顯式建立新列 PARTITIONED BY (days(o_orderdate)) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
高頻事件按小時分區
XIHE
-- 假設 l_receiptdate 擴充為 TIMESTAMP(類比收貨時間戳記) CREATE TABLE db_iceberg.lineitem_realtime ( l_orderkey BIGINT, l_partkey BIGINT, l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15,2), l_extendedprice DECIMAL(15,2), l_discount DECIMAL(15,2), l_tax DECIMAL(15,2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, l_commitdate DATE, l_receipttime TIMESTAMP, -- 類比:精確到秒的收貨時間 l_shipmode STRING ) PARTITIONED BY (hour(l_receipttime)) -- 按小時分區,支援近即時監控 STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.lineitem_realtime ( l_orderkey BIGINT, l_partkey BIGINT, l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15,2), l_extendedprice DECIMAL(15,2), l_discount DECIMAL(15,2), l_tax DECIMAL(15,2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, l_commitdate DATE, l_receipttime TIMESTAMP, l_shipmode STRING ) USING iceberg -- Iceberg 特性:使用 hours() 轉換函式進行隱藏式磁碟分割,無需新增列 PARTITIONED BY (hours(l_receipttime)) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
基於高基數欄位雜湊分桶
XIHE
-- 對高基數外鍵 l_partkey 雜湊分桶,避免小檔案 & 提升 JOIN 效能 CREATE TABLE db_iceberg.lineitem ( l_orderkey BIGINT, l_partkey BIGINT, -- 高基數(約 2000 萬唯一值) l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15,2), l_extendedprice DECIMAL(15,2), l_discount DECIMAL(15,2), l_tax DECIMAL(15,2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, l_commitdate DATE, l_receiptdate DATE, l_shipmode STRING ) PARTITIONED BY (bucket(l_partkey,64)) -- 分64桶,均勻分布 STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.lineitem ( l_orderkey BIGINT, l_partkey BIGINT, -- High Cardinality l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15,2), l_extendedprice DECIMAL(15,2), l_discount DECIMAL(15,2), l_tax DECIMAL(15,2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, l_commitdate DATE, l_receiptdate DATE, l_shipmode STRING ) USING iceberg -- Iceberg 文法:bucket(桶數, 列名) -- 這會產生一個隱式分區列,根據雜湊值將資料分散到 64 個邏輯桶中 PARTITIONED BY (bucket(64, l_partkey)) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
truncate[len]字串首碼截斷分區
XIHE
-- 按電話號碼國家代碼首碼分區(如 '13-' 代表中國某電訊廠商) CREATE TABLE db_iceberg.customer_by_phone ( c_custkey BIGINT, c_name STRING, c_phone STRING, -- 格式:'13-888-999-1234' c_acctbal DECIMAL(15,2), c_mktsegment STRING ) PARTITIONED BY (truncate(c_phone,3)) -- 截取前3字元 '13-' STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.customer_by_phone ( c_custkey BIGINT, c_name STRING, c_phone STRING, c_acctbal DECIMAL(15,2), c_mktsegment STRING ) USING iceberg -- Iceberg 原生支援 truncate(寬度, 列名) 轉換 -- 資料寫入時,Iceberg 會自動計算 c_phone 的前3位並據此存入對應目錄 PARTITIONED BY (truncate(3, c_phone)) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
多級複合分區
兩級分區在分區數量(避免過細)與資料均勻性(避免傾斜)之間取得平衡,適用于海量資料的查詢情境。
以通用日誌表為例,建立語句如下:
XIHE
-- 使用者行為日誌表:時間 + 高基數ID分桶 + 地區截斷,三層複合分區 CREATE TABLE db_iceberg.user_event_log ( event_id BIGINT, user_id BIGINT, -- 高基數使用者ID session_id STRING, event_type STRING, event_time TIMESTAMP, -- 精確到秒的時間戳記 country_code STRING, -- 國家代碼,如 'CN', 'US', 'DE' device_type STRING, payload STRING ) PARTITIONED BY ( day(event_time), -- 第一級:按天分區(高效時間裁剪) bucket(user_id,64), -- 第二級:對高基數 user_id 雜湊分64桶(防小檔案) truncate(country_code,2) -- 第三級:截取國家代碼前2位(地區彙總) ) STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';-- TPC-H lineitem 表:按發貨日期天分區 + 零件ID分桶(經典事實表最佳化) CREATE TABLE db_iceberg.lineitem_multiple_part ( l_orderkey BIGINT, l_partkey BIGINT, -- 高基數外鍵(~20M 唯一值) l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15, 2), l_extendedprice DECIMAL(15, 2), l_discount DECIMAL(15, 2), l_tax DECIMAL(15, 2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, -- TPC-H 核心時間欄位 l_commitdate DATE, l_receiptdate DATE, l_shipinstruct STRING, l_shipmode STRING, l_comment STRING ) PARTITIONED BY ( day(l_shipdate), -- 第一級:按發貨日期分區(TPC-H 查詢高頻過濾條件) bucket(l_partkey,32) -- 第二級:對零件ID雜湊分32桶(提升 JOIN 和點查效能) ) STORED AS ICEBERG LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';Spark SQL
CREATE TABLE db_iceberg.user_event_log ( event_id BIGINT, user_id BIGINT, session_id STRING, event_type STRING, event_time TIMESTAMP, country_code STRING, device_type STRING, payload STRING ) USING iceberg PARTITIONED BY ( days(event_time), -- 自動轉換:按天分區 bucket(64, user_id), -- 自動雜湊:對 user_id 分64個桶 truncate(2, country_code) -- 自動截取:前2位字元 ) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';CREATE TABLE db_iceberg.lineitem_multiple_part ( l_orderkey BIGINT, l_partkey BIGINT, l_suppkey BIGINT, l_linenumber INT, l_quantity DECIMAL(15, 2), l_extendedprice DECIMAL(15, 2), l_discount DECIMAL(15, 2), l_tax DECIMAL(15, 2), l_returnflag STRING, l_linestatus STRING, l_shipdate DATE, l_commitdate DATE, l_receiptdate DATE, l_shipinstruct STRING, l_shipmode STRING, l_comment STRING ) USING iceberg PARTITIONED BY ( days(l_shipdate), -- 顯式指定按天分區 bucket(32, l_partkey) -- 雜湊分32桶 (注意參數順序:桶數在前) ) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
在AnalyticDB for MySQL中託管湖儲存中建立表
建立一個關聯到託管儲存的外部庫。
XIHE
CREATE EXTERNAL DATABASE test_db WITH DBPROPERTIES ('adb_lake_bucket' = '<YOUR_ADB_BUCKET>');Spark SQL
CREATE DATABASE test_db WITH DBPROPERTIES ('adb_lake_bucket' = '<YOUR_ADB_BUCKET>');
在該資料庫下建立表,並通過
TBLPROPERTIES定義表所位於的管理的資料湖bucket。說明分區策略與在使用者自有OSS bucket中建立一致。
XIHE
CREATE TABLE test_db.test_iceberg_tbl ( `id` int, `name` string ) STORED AS ICEBERG TBLPROPERTIES ( 'catalog_type' = 'ADB', 'adb_lake_bucket' = '<YOUR_ADB_BUCKET>' );Spark SQL
Job型資源群組
SET spark.adb.lakehouse.enabled=true; -- 開啟湖儲存 CREATE TABLE test_db.test_iceberg_tbl ( `id` int, `name` string ) USING iceberg TBLPROPERTIES ( 'adb_lake_bucket' = '<YOUR_ADB_BUCKET>' );Interactive型資源群組
開啟湖儲存。修改資源群組,添加Spark配置
spark.adb.lakehouse.enabled,值為true。執行SQL。
CREATE TABLE test_db.test_iceberg_tbl ( `id` int, `name` string ) USING iceberg TBLPROPERTIES ( 'adb_lake_bucket' = '<YOUR_ADB_BUCKET>' );
建立Delta Lake表
目前僅支援通過Spark SQL、PySpark建立和讀寫Delta Lake表,不支援通過XIHE引擎建立和讀寫。
建立樣本如下,文法說明及更多資訊,請參見How to Create Delta Lake Tables | Delta Lake。
CREATE DATABASE db_delta LOCATION 'oss://<YOUR_BUCKET>/db_delta/';
CREATE TABLE IF NOT EXISTS db_delta.delta_lake_comprehensive_test (
transaction_id BIGINT NOT NULL COMMENT '全域唯一交易ID',
user_id STRING COMMENT '使用者ID',
device_info STRUCT<
os: STRING,
model: STRING,
app_version: STRING
> COMMENT '裝置詳情嵌套結構',
tags MAP<STRING, STRING> COMMENT '使用者標籤Map',
item_list ARRAY<STRING> COMMENT '購買商品列表',
event_ts TIMESTAMP COMMENT '事件發生時間',
revenue DECIMAL(18, 2) COMMENT '營收金額',
event_date DATE COMMENT '自動產生的日期分區鍵'
)
USING DELTA
PARTITIONED BY (event_date);