Distributed Spark execution
===========================

Postgres Analytics AcceleratorPGAAをApache
Sparkと統合して、PostgresをマルチテラバイトのデータセットのCPUボトルネックを克服する高性能プラットフォームに変換します。

PGAAを介してクエリ実行をApache
Sparkにオフロードすると、アーキテクチャの利点がすぐに得られます。

- **意図的なハンドオフ**
  実行エンジンを明示的に選択することにより、完全な制御を維持します。
  PGAAは次に、 SQLクエリーをパッケージ化して、Spark
  Connectエンドポイントにルーティングします。

- **分散処理**
  ワークロードをSparkクラスターに移動することにより、分散コンピューティング能力を活用して、単一のデータベースインスタンスをボトルネックにしてしまう大規模なデータセットを処理します。

- **ゼロインデックス**
  手動インデックス作成や定期的なメンテナンスを必要とせずにピークの分析パフォーマンスを達成し、操作オーバーヘッドを大幅に削減します。

..  Note::
   現在、Apache Sparkとの統合は、以下をサポートしています。-S3互換オブジェクトストレージまたは共有POSIXファイルシステムのParquetファイルの読み取り専用クエリ。  - Iceberg RESTカタログのIcebergテーブルの読み取り専用クエリー。 !!!

Sparkを使用したPGAAの構成は、次の手順で構成されます。

1. :ref:`Spark環境の構成 <Spark環境の構成>`  。

2. :ref:`Connecting Postgres to Spark <#connecting-postgres-to-spark>`  。

3. :ref:`ソースデータの構成 <ソースデータの構成>`  。

4. :ref:`分析クエリの実行 <分析クエリの実行>`  。

5. :ref:`クエリプランの検査 <クエリプランの検査>`  。

Spark環境の構成
---------------

3. 4以降を実行しているSparkクラスターインストールを利用またはインストールできます。

次の依存関係で構成した実行中の `Spark Connect <https://spark.apache.org/docs/3.5.7/api/python/getting_started/quickstart_connect.html>`_ サーバーが必要です。

::

     org.apache.spark:spark-connect_2.12:3.5.6,\
     io.delta:delta-spark_2.12:3.3.1,\
     org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.9.2,\
     org.apache.iceberg:iceberg-aws-bundle:1.9.2,\
     org.apache.hadoop:hadoop-aws:3.3.4

..  Important::
   Iceberg Sparkライブラリーバージョン1.9.2には、Sparkを使用してIcebergテーブルから読み取る場合に既知の問題があります。特定の条件下では、同時実行中に等価削除がスキップされる場合があります。現在の回避策は、`spark-defaults.conf` ファイルでSparkアプリケーション設定 `spark.sql.iceberg.executor-cache.enabled <https://iceberg.apache.org/docs/latest/spark-configuration/>`_ を無効にすることです。

このキャッシュを無効にすると、すべての削除が正しく処理されるため、データの一貫性が保証されますが、大容量の読み取りワークロードの場合、パフォーマンスに影響を与える可能性があります。

.. ::
   ## PostgresをSparkに接続する

1. 実行エンジンを設定し、PostgresセッションでSpark Connect
   URLを定義します。

.. code:: sql

     SET pgaa.executor_engine = spark_connect;
     SET pgaa.spark_connect_url = sc://spark-connect:15002;

``spark-connect`` は、Spark Connectサービスアドレスを指します。

1. PGAAインターフェイスを介して単純なバージョンチェックを実行することにより、接続を検証します。

.. code:: sql

     SELECT pgaa.spark_sql(SELECT version());

成功した場合、コマンドはSparkクラスターのバージョン文字列を結果ます。

ソースデータの構成
------------------

PGAAは、データの保存場所としてS3互換ストレージおよびローカルファイルシステムからのParquetファイルの読み取り、またはIceberg
RESTカタログのIcebergテーブルの読み取りをサポートしています。

オプション1 Parquetファイルからの読み取り
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

1. Postgresデータベースに、分析データが含まれるバケットを指す 
`PGFS storage location <https://www.enterprisedb.com/docs/edb-postgres-ai/latest/ai-factory/pipeline/pgfs/functions/#creating-a-storage-location>`_ を作成します。例、パブリックバケットの場合

.. code:: sql

     SELECT pgfs.create_storage_location(
     my-sample-data,
     s3://beacon-analytics-demo-data-us-east-1-prod,
     {"skip_signature": "true", "region": "us-east-1"}
     );

1. ``PGAA``
   アクセス方法を使用してテーブルを作成します。前の手順で作成した保存場所、単一のParquetファイルまたは複数のParquetファイルを含むディレクトリへのパス、およびフォーマット現在Parquetを指定します。例

.. code:: sql

     CREATE TABLE customer () USING PGAA WITH (pgaa.storage_location = my-sample-data, pgaa.path = tpch_sf_1/customer, pgaa.format = parquet);
     CREATE TABLE orders () USING PGAA WITH (pgaa.storage_location = my-sample-data, pgaa.path = tpch_sf_1/orders, pgaa.format = parquet);
     CREATE TABLE lineitem () USING PGAA WITH (pgaa.storage_location = my-sample-data, pgaa.path = tpch_sf_1/lineitem, pgaa.format = parquet);
     CREATE TABLE nation () USING PGAA WITH (pgaa.storage_location = my-sample-data, pgaa.path = tpch_sf_1/nation, pgaa.format = parquet);

複数のParquetファイルを含むディレクトリを指定すると、PGAAはすべてのファイルを処理のために単一のテーブルに自動的に結合します。

オプション2 Icebergテーブルからの読み取り
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

1. カタログ接続を構成します。例

.. code:: sql

     SELECT pgaa.add_catalog(
       my_iceberg_catalog,
       iceberg-rest,
       {
         "url": "https://iceberg.example.com/api/v1",
         "warehouse": "a1b2c3d4-e5f6-g7h8",
         "token": "secret_token_abc123"
       }
     );

1. 必要なカタログ管理テーブルを作成します。または、カタログを添付またはインポートします。

.. code:: sql

     CREATE TABLE customer () USING PGAA WITH (pgaa.format = iceberg,pgaa.managed_by = my_iceberg_catalog,pgaa.catalog_namespace = public,pgaa.catalog_table = customer);
     CREATE TABLE orders () USING PGAA WITH (pgaa.format = iceberg,pgaa.managed_by = my_iceberg_catalog,pgaa.catalog_namespace = public,pgaa.catalog_table = orders);
     CREATE TABLE lineitem () USING PGAA WITH (pgaa.format = iceberg,pgaa.managed_by = my_iceberg_catalog,pgaa.catalog_namespace = public,pgaa.catalog_table = lineitem);
     CREATE TABLE nation () USING PGAA WITH (pgaa.format = iceberg,pgaa.managed_by = my_iceberg_catalog,pgaa.catalog_namespace = public,pgaa.catalog_table = nation);

Icebergカタログを使用してテーブルへのアクセスを構成する方法の詳細については、
:ref:`Integrating with Iceberg catalogs <Integrating with Iceberg catalogs>` を参照してください。

分析クエリの実行
----------------

この例のクエリーは、特定の四半期内に最高の収益を達成した顧客を特定することにより、トップ20収益レポートを生成します。これをSparkにオフロードすることにより、PGAAは複雑な条件付き集計を並列化し、述語プッシュダウンを利用してオブジェクトストレージからのデータ転送を最小限に抑えます。

.. code:: sql

   SELECT
       c_custkey,
       c_name,
       SUM(l_extendedprice * (1 - l_discount)) AS revenue,
       c_acctbal,
       n_name,
       c_address,
       c_phone,
       c_comment
   FROM
       customer,
       orders,
       lineitem,
       nation
   WHERE
       c_custkey = o_custkey
       AND l_orderkey = o_orderkey
       AND o_orderdate >= 1993-10-01
       AND o_orderdate < 1994-01-01
       AND l_returnflag = R
       AND c_nationkey = n_nationkey
   GROUP BY
       c_custkey,
       c_name,
       c_acctbal,
       c_phone,
       n_name,
       c_address,
       c_comment
   ORDER BY
       revenue DESC,
       c_custkey
   LIMIT 20;

クエリプランの検査
------------------

Sparkが有効になっていない場合、PGAAはローカルベクトル化実行のためにデフォルトのSeafowlエンジンを使用します。
``EXPLAIN``
をステートメントの前に追加して、ローカルのシングルノード実行からSpark
Connectを介したリモートの分散処理モデルへの移行を確認します。

これらの2つのクエリプランを比較すると、ローカルのシングルノード実行からSpark
Connectを介したリモートの分散処理モデルへの移行が強調表示されます。

- **結合戦略**
  Seafowlは、データベースに行をシークエンシャルに比較することを強制する「ネストループ」アプローチ
  ``Cross Join`` + ``Filter`` を使用します。
  Sparkはこれを\ ``BroadcastHashJoin``
  にアップグレードします。ここでは、小さなテーブルがすべてのワーカーノードに送信されて、大きなファクトテーブルを並列処理します。

- **集計ロジック** シングルノードシークエンシャル\ ``aggregate``
  の代わりに、Sparkは2ステージ\ ``HashAggregate``
  を使用します。これは、個々のワーカーノードで\ ``partial_sum``
  を計算し、\ ``Exchange``
  シャッフルを実行して結果を結合し、大規模データのホリゾンタル
  スケーラビリティを提供します。

- **データアクセス**
  Seafowlは、使用可能なすべてのデータに対して標準の\ ``TableScan``
  を実行します。 Sparkは、Iceberg
  RESTカタログを介してメタデータを認識する\ ``BatchScan``
  を使用します。これは、Icebergマニフェストファイルを読み取り、S3から取得される前に無関係なデータブロックを完全にスキップします。

Sparkプランは、大規模なデータセットのクラスター全体で水平にスケーリングするように設計された洗練された分散ブループリントです。逆に、Seafowlプランは、ネットワークレイテンシーとクラスター調整のオーバーヘッドが実際の処理時間を上回る場合に、小規模なデータセット用に最適化された高速のローカライズされたブループリントです。
