Distributed Spark execution#

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

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

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

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

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

!!!注

現在、Apache Sparkとの統合は以下をサポートしています。

  • S3互換オブジェクトストレージまたは共有POSIXファイルシステムのParquetファイルの読み取り専用クエリー。

  • Iceberg RESTカタログのIcebergテーブルの読み取り専用クエリーOAuth2認証はサポートされていません。

  1. Spark環境の構成

  2. PostgresをSparkに接続する

  3. ソースデータの構成

  4. 分析クエリの実行

  5. クエリプランの検査

Spark環境の構成#

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

次の依存関係で構成した実行中の Spark Connect サーバーが必要です。

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

!!!重要

Iceberg Sparkライブラリーバージョン1.9.2には、Sparkを使用してIcebergテーブルから読み取る場合に既知の問題があります。特定の条件下では、同時実行中に等価削除がスキップされる場合があります。

現在の回避策は、spark-defaults.conf ファイルでSparkアプリケーション設定 spark.sql.iceberg.executor-cache.enabled を無効にすることです。

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

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

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

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

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

SELECT pgaa.spark_sql(SELECT version());

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

ソースデータの構成#

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

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

  1. Postgresデータベースに、分析データが含まれるバケットを指す PGFS storage location を作成します。例、パブリックバケットの場合

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を指定します。例

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. カタログ接続を構成します。例

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. 必要なカタログ管理テーブルを作成します。または、カタログを添付またはインポートします。

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

分析クエリの実行#

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

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プランは、ネットワークレイテンシーとクラスター調整のオーバーヘッドが実際の処理時間を上回る場合に、小規模なデータセット用に最適化された高速のローカライズされたブループリントです。