Reading from object storage on WarehousePG#
WarehousePGのPGAAは、オブジェクトストレージのデータにアクセスする3つの方法をサポートしています。各オプションは独立しています。セットアップに合ったものを使用してください。
ストレージの場所を使用したオブジェクトストレージの照会#
ストレージ場所はオブジェクトストレージに直接接続し、Parquet、Delta Lake、およびIcebergファイルをサポートします。コーディネーターで保存場所とPGAAテーブルを作成し、標準のSQLを使用してデータを照会します。 WarehousePGは、最も効率の良い実行モードを自動的に選択し、クエリーの性質に応じて、コーディネーターでクエリーを実行またはセグメントホスト全体に分散します。
!!!注
WarehousePGのPGAAテーブルは読み取り専用です。データはオブジェクトストレージに既に存在している必要があります。
CREATE TABLE AS SELECT
CTASを使用したオブジェクトストレージへのデータの書き込みは、サポートされていません。
保存場所の構成の詳細については、 Reading tables in object storage を参照してください。
例#
次の例では、TPC-H SF1データを含むパブリックS3バケットに接続し、3つのテーブルを登録し、分析クエリを実行して、DirectScan純粋なPGAAとCompatScanネイティブヒープテーブルと結合したPGAAの両方を示します。
コーディネーターで保存場所を作成します。
SELECT pgfs.create_storage_location(
sample-data,
s3://beacon-analytics-demo-data-us-east-1-prod,
{"aws_skip_signature": "true"}
);
データを指すPGAAテーブルを作成します。
CREATE TABLE customer () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/customer);
CREATE TABLE orders () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/orders);
CREATE TABLE lineitem () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/lineitem);
分析クエリを実行します。
SELECT
l_orderkey,
sum(l_extendedprice * (1 - l_discount)) AS revenue,
o_orderdate,
o_shippriority
FROM
customer,
orders,
lineitem
WHERE
c_mktsegment = BUILDING
AND c_custkey = o_custkey
AND l_orderkey = o_orderkey
AND o_orderdate < DATE 1995-03-15
AND l_shipdate > DATE 1995-03-15
GROUP BY
l_orderkey,
o_orderdate,
o_shippriority
ORDER BY
revenue DESC,
o_orderdate
LIMIT 10;
EXPLAINを使用して、クエリプランを観察します。 PGAAは最初にDirectScanを試行し、コーディネーター上のSeafowlに実行を完全にオフロードします。このプランには、述語プッシュダウン、集計、ソートを含む完全なSeafowl論理プランが表示されます。 WarehousePGプランナーが関係していないため、Optimizer行が表示されないことに注意してください。
QUERY PLAN
--------------------------------------------------------------------------------------------
SeafowlDirectScan: Logical Plan
Limit: skip=0, fetch=10
Sort: revenue DESC NULLS FIRST, orders.o_orderdate ASC NULLS LAST
Projection: lineitem.l_orderkey, sum(lineitem.l_extendedprice * Int64(1) - lineitem.l_discount) AS revenue, orders.o_orderdate, orders.o_shippriority
Aggregate: groupBy=[[lineitem.l_orderkey, orders.o_orderdate, orders.o_shippriority]], aggr=[[sum(lineitem.l_extendedprice * (Int64(1) - lineitem.l_discount))]]
Filter: customer.c_mktsegment = Utf8("BUILDING") AND customer.c_custkey = orders.o_custkey AND lineitem.l_orderkey = orders.o_orderkey AND orders.o_orderdate < CAST(Utf8("1995-03-15") AS Date32) AND lineitem.l_shipdate > CAST(Utf8("1995-03-15") AS Date32)
Cross Join:
Cross Join:
TableScan: customer
TableScan: orders
TableScan: lineitem
(11 rows)
クエリにネイティブヒープテーブルとの結合など、Seafowlが処理できない操作が含まれる場合、PGAAはCompatScanにフォールバックします。例、
customer とネイティブsegments テーブルを結合する場合
CREATE TABLE segments (seg_id INT, seg_name TEXT) DISTRIBUTED BY (seg_id);
INSERT INTO segments VALUES (1, BUILDING), (2, MACHINERY);
EXPLAIN SELECT c_name, seg_name
FROM customer
JOIN segments ON c_mktsegment = seg_name;
PGAAは最初にDirectScanを試行し、次の通知を使用してCompatScanにフォールバックします。
NOTICE: PGAA DirectScan failed: QueryPlanningError("Error during planning: table default.public.segments not found"). Running the query in compatibility mode.
QUERY PLAN
--------------------------------------------------------------------------------------------------------
Gather Motion 4:1 (slice1; segments: 4) (cost=1414.00..246348.40 rows=748960 width=64)
-> Hash Join (cost=1414.00..236986.40 rows=187240 width=64)
Hash Cond: (customer.c_mktsegment = segments.seg_name)
-> Custom Scan (SeafowlCompatScan) on customer (cost=25.00..1131.25 rows=37500 width=64)
SeafowlPlan: Logical Plan
Projection: r1.c_name, r1.c_mktsegment
SubqueryAlias: r1
TableScan: public.customer
SeafowlQuery: select r1.c_name, r1.c_mktsegment from public.customer as r1
-> Hash (cost=769.00..769.00 rows=49600 width=32)
-> Broadcast Motion 4:4 (slice2; segments: 4) (cost=0.00..769.00 rows=49600 width=32)
-> Seq Scan on segments (cost=0.00..149.00 rows=12400 width=32)
Optimizer: Postgres-based planner
(13 rows)
Gather Motion 4:1 は、Strewnローカスプランを示します。
SeafowlCompatScan
は、4つのセグメントホストの各ホストで並列実行され、それぞれがcustomer
データのスライスを読み取り、結果がコーディネーターに収集されます。ヒープテーブルとの結合は、WarehousePGプランナーによって処理されます。下部のOptimizer: Postgres-based planner
行は、Seafowlがすべてを独立して処理するDirectScanとは異なり、Postgresプランナーがプラン全体を担当したことを確認します。
PGAAがGeneral遺伝子座とStrewn遺伝子座を選択する方法の詳細については、 クエリ実行モード を参照してください。
Icebergカタログに接続する#
WarehousePGコーディネーターに外部Iceberg RESTカタログを接続して、クラスター全体でカタログ管理テーブルを照会可能にします。
WarehousePGでPGAAを使用して、 EDB Postgres DistributedPGDがオブジェクトストレージにレプリケートしたデータを読み取ることができます。 PGDは、Iceberg形式のデータを書き込み、Iceberg RESTカタログに登録します。WarehousePGは、データを移動せずにMPPスケールで接続および照会できます。この設定のPGD側については、 PGDを使用した複製 を参照してください。
!!!注 pgaa.attach_catalog()
を使用した自動カタログ同期は、WarehousePGではサポートされていません。カタログが変更されるたびに、
pgaa.import_catalog() を手動で実行する必要があります。
資格情報構成やサポートされているカタログタイプを含む完全なセットアップ手順については、 Integrating with Iceberg catalogs を参照してください。
例#
コーディネーターにカタログを登録します。
SELECT pgaa.add_catalog(
my_catalog,
iceberg-rest,
{"uri": "https://iceberg-catalog.example.com"}
);
テーブルのメタデータをインポートします。
SELECT pgaa.import_catalog(my_catalog);
オプションで、インポートを特定の名前空間に制限します。
SELECT pgaa.import_catalog(my_catalog, my_schema);
インポートされると、テーブルはWarehousePGセッションから照会できます。
例 で説明されているのと同じクエリプランパターンが適用されます。
EXPLAIN
を使用して、PGAAがDirectScanまたはCompatScanを使用しているか、およびどの遺伝子座が選択されているかを観察します。
Sparkでの加速#
Spark Connectを介してクエリーを外部Apache Sparkクラスターにルーティングします。このモードでは、コーディネーターはクエリーをSparkに直接転送し、セグメントホストは実行に関与せず、SparkはWarehousePGプランナーとは無関係に独自の計画と配布を処理します。
完全なセットアップ手順、GPUアクセラレーションオプション、およびパフォーマンスガイダンスについては、 Accelerating with Spark を参照してください。
例#
pgaa.spark_connect_urlを設定し、コーディネーターの実行エンジンを切り替えます。
SET pgaa.executor_engine = spark_connect;
SET pgaa.spark_connect_url = sc://spark-connect-hostname:15002;
例 と同じパブリックTPC-Hサンプルデータを使用して、保存場所とPGAAテーブルを作成します。 Spark ConnectはParquetのみをサポートしているため、
pgaa.formatを明示的に設定する必要があります。
SELECT pgfs.create_storage_location(
sample-data,
s3://beacon-analytics-demo-data-us-east-1-prod,
{"aws_skip_signature": "true"}
);
CREATE TABLE customer () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/customer, pgaa.format = parquet);
CREATE TABLE orders () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/orders, pgaa.format = parquet);
CREATE TABLE lineitem () USING PGAA
WITH (pgaa.storage_location = sample-data, pgaa.path = tpch_sf_1/lineitem, pgaa.format = parquet);
同じTPC-H Q3クエリーを実行します。コーディネーターは実行をSparkに転送します。
SELECT
l_orderkey,
sum(l_extendedprice * (1 - l_discount)) AS revenue,
o_orderdate,
o_shippriority
FROM
customer,
orders,
lineitem
WHERE
c_mktsegment = BUILDING
AND c_custkey = o_custkey
AND l_orderkey = o_orderkey
AND o_orderdate < DATE 1995-03-15
AND l_shipdate > DATE 1995-03-15
GROUP BY
l_orderkey,
o_orderdate,
o_shippriority
ORDER BY
revenue DESC,
o_orderdate
LIMIT 10;
EXPLAINを使用して、クエリプランを観察します。 Sparkがアクティブなエンジンの場合、PGAAはSparkDirectScanノードを生成します。このプランは、Sparkクラスター内で完全に処理される述語プッシュダウン、結合、集計を含むSpark独自の実行を示します。Optimizer行は表示されず、 WarehousePGプランナーが関与していないことを確認します。
QUERY PLAN
--------------------------------------------------------------------------------------------
SparkDirectScan: Logical Plan
AdaptiveSparkPlan: isFinalPlan=false
+- TakeOrderedAndProject: (limit=10, orderBy=[revenue DESC, o_orderdate ASC])
+- HashAggregate: (keys=[l_orderkey, o_orderdate, o_shippriority], functions=[sum(...)])
+- HashAggregate: (keys=[l_orderkey, o_orderdate, o_shippriority], functions=[partial_sum(...)])
+- Project: [o_orderdate, o_shippriority, l_orderkey, l_extendedprice, l_discount]
+- SortMergeJoin: [o_orderkey], [l_orderkey], Inner
:- Sort [o_orderkey ASC NULLS FIRST]
: +- Exchange hashpartitioning(o_orderkey, 200)
: +- BroadcastHashJoin [c_custkey], [o_custkey], Inner, BuildLeft
: :- BroadcastExchange
: : +- Filter (c_mktsegment = BUILDING)
: : +- FileScan parquet [c_custkey, c_mktsegment]
: : Location: s3a://beacon-analytics-demo-data-us-east-1-prod/tpch_sf_1/customer
: +- Filter (o_orderdate < 1995-03-15)
: +- FileScan parquet [o_orderkey, o_custkey, o_orderdate, o_shippriority]
: Location: s3a://beacon-analytics-demo-data-us-east-1-prod/tpch_sf_1/orders
+- Sort [l_orderkey ASC NULLS FIRST]
+- Exchange hashpartitioning(l_orderkey, 200)
+- Filter (l_shipdate > 1995-03-15)
+- FileScan parquet [l_orderkey, l_extendedprice, l_discount, l_shipdate]
Location: s3a://beacon-analytics-demo-data-us-east-1-prod/tpch_sf_1/lineitem
(22 rows)
完全なSpark実行プランを検査するには、ブラウザーで
http://<spark-host>:4040のSpark UIを開きます。 [ SQL ]タブに移動して、ステージ境界、シャッフル操作、ステージごとのメトリックを含むクエリプランの視覚的な内訳を表示します。