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の両方を示します。

  1. コーディネーターで保存場所を作成します。

SELECT pgfs.create_storage_location(
  sample-data,
  s3://beacon-analytics-demo-data-us-east-1-prod,
  {"aws_skip_signature": "true"}
);
  1. データを指す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);
  1. 分析クエリを実行します。

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;
  1. 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 を参照してください。

例#

  1. コーディネーターにカタログを登録します。

SELECT pgaa.add_catalog(
  my_catalog,
  iceberg-rest,
  {"uri": "https://iceberg-catalog.example.com"}
);
  1. テーブルのメタデータをインポートします。

SELECT pgaa.import_catalog(my_catalog);
  1. オプションで、インポートを特定の名前空間に制限します。

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 を参照してください。

例#

  1. pgaa.spark_connect_url を設定し、コーディネーターの実行エンジンを切り替えます。

SET pgaa.executor_engine = spark_connect;
SET pgaa.spark_connect_url = sc://spark-connect-hostname:15002;
  1. 例 と同じパブリック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);
  1. 同じ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;
  1. 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)
  1. 完全なSpark実行プランを検査するには、ブラウザーでhttp://<spark-host>:4040 のSpark UIを開きます。 [ SQL ]タブに移動して、ステージ境界、シャッフル操作、ステージごとのメトリックを含むクエリプランの視覚的な内訳を表示します。