AnalyticDB for MySQL は、Apache Iceberg などのオープンソースのレイクテーブルフォーマットと深く統合されています。自社開発の高性能 XIHE エンジンとマネージド Spark エンジンを搭載し、オープンでマルチエンジン互換のデータレイク (レイクハウス) 機能を提供します。レイクテーブルのデータは、標準の Parquet フォーマットでオブジェクトストレージ (Alibaba Cloud OSS) に永続化されます。Spark、Flink、Trino など、Iceberg をサポートするあらゆるコンピュートエンジンがデータを直接読み取ることができます。これにより、ベンダーロックインを回避し、データ資産の長期的な可用性を確保できます。このトピックでは、Apache Iceberg と Delta Lake フォーマットを例に、データレイクテーブルの作成方法について説明します。
ストレージモード
モードの選択
データレイクテーブルを作成する前に、データの保存場所を決定する必要があります。AnalyticDB for MySQL は、コントロールと利便性のバランスを取り、さまざまなセキュリティおよび O&M 要件を満たす、2 つの柔軟なストレージ管理モードを提供します。
お客様が管理する OSS バケット
データは、お客様が指定した同一リージョン内の Alibaba Cloud OSS バケットに完全に保存されます。このアプローチは、厳格なコンプライアンスとデータ主権の要件を満たします。きめ細かな制御を行うには、データベースとテーブルを作成する際にストレージパスを明示的に宣言する必要があります。
AnalyticDB for MySQL が管理するレイクストレージ
AnalyticDB for MySQL は、基盤となるストレージバケットを自動的に管理します。これらのバケットは、お客様のアカウントには表示されません。標準 SQL を使用して、ファイルシステム、権限設定、またはライフサイクル管理を行うことなく、レイクテーブルの読み書きをシームレスに行うことができます。これにより、導入のハードルが大幅に下がります。詳細については、「レイクストレージ」をご参照ください。
主な利点
データベース、テーブル構造、列定義、パーティション情報など、レイクテーブルのすべてのメタデータは、AnalyticDB for MySQL の組み込みカタログサービスによって一元管理されます。個別のメタデータクラスターをデプロイ、スケーリング、または維持する必要はありません。
オープンデータレイクの操作は、従来のデータベースを使用するのと同じくらい簡単です。
CREATE TABLEと SQL クエリロジックにのみ集中すればよいです。プラットフォームは、ストレージ、メタデータ、ファイル形式、圧縮コーデックなど、基盤となるインフラストラクチャを完全に管理します。
オープンソースエコシステムへの完全なオープン性が維持され、「データベースのようにシンプルで、データレイクのようにオープン」という理想的なレイクハウス体験を提供します。
前提条件
XIHE エンジンを使用して Apache Iceberg テーブルを作成するには、ご利用のクラスターのマイナーエンジンバージョンが 3.2.7 以降である必要があります。
説明マイナーバージョンの表示と更新を行うには、AnalyticDB for MySQL コンソールの クラスター情報 ページにある 構成情報 セクションに移動します。すでに最新のデフォルトベースラインバージョンのクラスターをアップグレードするには、DingTalk (DingTalk ID:
x5v_rm8wqzuqf) で Alibaba Cloud サービスサポートにご連絡ください。AnalyticDB for MySQL のマネージドレイクストレージにテーブルを作成するには、チケットを送信してテクニカルサポートに連絡し、レイクストレージ機能を有効にして、新しいレイクストレージを作成する必要があります。
デフォルトのテーブルフォーマット
データベースレベルでデフォルトのテーブルフォーマットを設定できます。これにより、データベースで作成される新しいテーブルは、各 CREATE TABLE 文で明示的に宣言しなくても、指定されたフォーマットを自動的に使用するようになります。
ジョブリソースグループ、対話型リソースグループ、または送信されたジョブで、次のパラメーターを設定します。
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 バケット内
非パーティション化テーブル
ユースケース:国、地域、製品カテゴリなどの小規模なディメンションテーブル、頻繁にスキャンされる小規模な静的データセット、または時間やカーディナリティの高いフィールドによるプルーニングを必要としないデータに最適です。
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 は、パーティション変換関数を使用してパーティションを定義することをサポートしています。パーティションテーブルは、高性能でスケーラブル、かつ管理しやすい最新のデータレイク (レイクハウス) を構築するための中心的な手法です。AnalyticDB for MySQL は、次の Apache Iceberg パーティション変換ルールをサポートしています。
変換関数 | 構文例 | 説明 | タイプ |
|
| ID パーティション分割。Hive スタイルのパーティション分割に相当します。 | すべてのタイプ。ただし、カーディナリティの高いフィールドには推奨されません。 |
|
| 年でパーティション分割します。 | Timestamp、Date |
|
| 月でパーティション分割します。 | Timestamp、Date |
|
| 日でパーティション分割します。これは最も一般的なアプローチです。 | Timestamp、Date |
|
| 時間でパーティション分割します。 | Timestamp |
|
| ハッシュバケット化。N はバケット数です。 | すべてのタイプ。この関数は、ID などのカーディナリティの高いフィールドによく使用されます。 |
|
| 文字列の最初の | String |
以下の例は、XIHE エンジンと Spark SQL を使用して、さまざまなユースケースにパーティション分割戦略を適用する方法を示しています。
カテゴリフィールド
XIHE
-- c_mktsegment (市場セグメント) でパーティション分割します。 -- このフィールドには 'AUTOMOBILE'、'BUILDING'、'FURNITURE' などの 5 つの一意の値しかなく、ID パーティション分割に最適です。 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 ) 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 -- 注:Spark の Iceberg は 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, -- 高カーディナリティ (約 2,000 万の一意の値) 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, -- 高カーディナリティ 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(width, column_name) 変換をネイティブにサポートしています。 -- データが書き込まれると、Iceberg は c_phone の最初の 3 文字を自動的に計算し、対応するディレクトリにデータを保存します。 PARTITIONED BY (truncate(3, c_phone)) LOCATION 'oss://<YOUR_OSS_BUCKET>/PATH/SUBPATH/';
複合パーティション分割
複数レベルのパーティション分割は、パーティション数のバランスを取り、過剰なパーティション分割を回避すると同時に、均等なデータ分散を確保してデータスキューを防ぎます。この戦略は、大規模なデータセットのクエリに最適です。
次の例は、汎用的なログテーブルの作成方法を示しています。
XIHE
-- ユーザーイベントログテーブル:時間、高カーディナリティ ID のバケット化、およびリージョン切り捨てを使用する 3 レベルの複合パーティション戦略。 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), -- レベル 1:日単位でパーティション分割し、効率的な時間ベースのプルーニングを実現します。 bucket(user_id,64), -- レベル 2:高カーディナリティの user_id に 64 バケットのハッシュバケット化を適用して、小さなファイルを防ぎます。 truncate(country_code,2) -- レベル 3:国コードの最初の 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, -- 高カーディナリティの外部キー (約 2,000 万の一意の値) 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_shipSpark 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を使用して、テーブルが保存されるマネージドデータレイクバケットを指定します。説明パーティション分割戦略は、お客様自身の Object Storage Service (OSS) バケットに作成されたテーブルの戦略と同じです。
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
ジョブリソースグループ
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>' );対話型リソースグループ
レイクストレージを有効にします。リソースグループを変更し、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 テーブルの作成
現在、Delta Lake テーブルの作成、読み取り、書き込みをサポートしているのは Spark SQL と PySpark のみです。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 'ユーザータグのマップ',
item_list ARRAY<STRING> COMMENT '購入済みアイテムのリスト',
event_ts TIMESTAMP COMMENT 'イベントが発生した時刻',
revenue DECIMAL(18, 2) COMMENT '収益額',
event_date DATE COMMENT '自動生成された日付パーティションキー'
)
USING DELTA
PARTITIONED BY (event_date);