PG-Strom v6.0特集:Arrow_Fdw仮想列

PG-Strom v6.0をリリースしました。

GPU-Sortと一部のWindow関数対応
マルチGPUのPinned Inner Buffer対応
Arrow_Fdwの仮想列機能
GPUでの完全な集計値の生成
といった、いくつかの重要機能を含むリリースで、特にGPU-Sortによって新しいワークロードへの対応が広がったという事でバージョン6.x系列としました。
その主要機能について、何回かに分けて解説していきたいと思います。

Arrow_Fdwで時系列のデータを管理する

トランザクショナルなデータを集計・分析する時に、こんな感じのデータ管理を考えた事はないでしょうか?

  • 更新・削除の可能性がある直近のデータはPostgreSQLのテーブル(Heapフォーマット/行形式)に保管しておこう。
  • 一定期間が過ぎたら、もう更新・削除が行われる事もないので、Arrow_Fdw(Arrowフォーマット/列形式)に変換しよう。
  • その時、Arrowファイルはパッと見で分類できるように、YYYYMMDDで日付や、店舗、地域などでカテゴリ分けしよう。

マメな人なら、ファイル名の一部が何らかのコンテンツの属性を代表しているという事もそう珍しい事ではないでしょう。

では、例えばこのような感じのArrowファイルを考えてみる事にします。

$ ls /opt/arrow/mydata/
f_lineorder_1993_AIR.arrow    f_lineorder_1995_RAIL.arrow
f_lineorder_1993_FOB.arrow    f_lineorder_1995_SHIP.arrow
f_lineorder_1993_MAIL.arrow   f_lineorder_1995_TRUCK.arrow
f_lineorder_1993_RAIL.arrow   f_lineorder_1996_AIR.arrow
f_lineorder_1993_SHIP.arrow   f_lineorder_1996_FOB.arrow
f_lineorder_1993_TRUCK.arrow  f_lineorder_1996_MAIL.arrow
f_lineorder_1994_AIR.arrow    f_lineorder_1996_RAIL.arrow
f_lineorder_1994_FOB.arrow    f_lineorder_1996_SHIP.arrow
f_lineorder_1994_MAIL.arrow   f_lineorder_1996_TRUCK.arrow
f_lineorder_1994_RAIL.arrow   f_lineorder_1997_AIR.arrow
f_lineorder_1994_SHIP.arrow   f_lineorder_1997_FOB.arrow
f_lineorder_1994_TRUCK.arrow  f_lineorder_1997_MAIL.arrow
f_lineorder_1995_AIR.arrow    f_lineorder_1997_RAIL.arrow
f_lineorder_1995_FOB.arrow    f_lineorder_1997_SHIP.arrow
f_lineorder_1995_MAIL.arrow   f_lineorder_1997_TRUCK.arrow

これは、Star Schema Benchmarkのlineorderテーブルを、そのタイムスタンプ(lo_orderdate列)とカテゴリ(lo_shipmode)毎に分類し、別々のファイルに保存したものです。
シェルスクリプトで記述すると以下のようになります。

for y in 1993 1994 1995 1996 1997 1998;
do
  for m in AIR RAIL FOB SHIP MAIL TRUCK RAIL;
  do
    pg2arrow -d ssbm -c "SELECT * FROM lineorder WHERE lo_orderdate BETWEEN ${y}0101 AND ${y}1231 AND lo_shipmode = '$m'" -o /opt/arrow/f_lineorder_${y}_${m}.arrow;
  done;
done

計算機ではなく人間がこれらのファイルを参照してデータを検索するとした場合、例えば『1996年3月のトラック輸送、または鉄道輸送の記録・・・』と検索のオーダーに対して、f_lineorder_1993_SHIP.arrowf_lineorder_1997_AIR.arrowのファイルをわざわざ探そうとするでしょうか?しませんよね😅

PG-Strom v6.0で導入されたArrow_Fdwの仮想列機能は、そうしたデータ管理のHeuristicsを一部の検索・集計にも付け加えてやろうというものです。

Arrow_Fdw仮想列の定義

こういったファイル名の命名規則は、サイトによっても異なりますし、管理者によっても異なります。
そこで、Arrow_Fdwはファイル名のどの部分に意味ある情報が含まれているかをオプションによって指定する事を可能にしました。

postgres=# IMPORT FOREIGN SCHEMA f_lineorder FROM SERVER arrow_fdw INTO public
           OPTIONS (dir '/opt/arrow/mydata', pattern 'f_lineorder_@{year}_${shipmode}.arrow');
IMPORT FOREIGN SCHEMA
postgres=# \d f_lineorder
                             Foreign table "public.f_lineorder"
       Column       |     Type      | Collation | Nullable | Default |     FDW options
--------------------+---------------+-----------+----------+---------+----------------------
 lo_orderkey        | numeric       |           |          |         |
 lo_linenumber      | integer       |           |          |         |
 lo_custkey         | numeric       |           |          |         |
 lo_partkey         | integer       |           |          |         |
 lo_suppkey         | numeric       |           |          |         |
 lo_orderdate       | integer       |           |          |         |
 lo_orderpriority   | character(15) |           |          |         |
 lo_shippriority    | character(1)  |           |          |         |
 lo_quantity        | numeric       |           |          |         |
 lo_extendedprice   | numeric       |           |          |         |
 lo_ordertotalprice | numeric       |           |          |         |
 lo_discount        | numeric       |           |          |         |
 lo_revenue         | numeric       |           |          |         |
 lo_supplycost      | numeric       |           |          |         |
 lo_tax             | numeric       |           |          |         |
 lo_commit_date     | character(8)  |           |          |         |
 lo_shipmode        | character(10) |           |          |         |
 year               | bigint        |           |          |         | (virtual 'year')
 shipmode           | text          |           |          |         | (virtual 'shipmode')
Server: arrow_fdw
FDW options: (dir '/opt/arrow/mydata', pattern 'f_lineorder_@{year}_${shipmode}.arrow')

PG-Strom v6.0では新たにArrow_Fdwのオプションとしてpatternが追加されました。
このオプションにより、dirで指定したディレクトリに存在するArrowファイル(複数可)をマッピングした外部テーブルを定義する際に、patternで指定したファイル名だけしか対象に含まれなくなります。
patternにはワイルドカードを含める事ができ、@{xxxx}${xxxx}で囲まれた部分はワイルドカードとして解釈されます。つまり、マッチングした部分に名前が付いていると理解してください。

IMPORT FOREIGN SCHEMAでこれらのArrowファイルをインポートすると、元となったlineorder表には含まれていない列が二つほど追加されているのが見えます。
year列は、ファイル名の@{year}部分にマッチした内容を示すもので、整数値に変換できることを前提にbigint型のデータとしてマップされます。
shipmode列は、ファイル名の${shipmode}部分にマッチした内容を示すもので、単純にマッチした部分をtext型のデータとしてマップされます。

ここで定義したf_lineorder表は30個のArrowファイルを含んでいるため、何も考えずにスキャンすれば30個のArrowファイルを順にスキャンする事になりますが、ある特定のArrowファイルを読み出している間はこれらの仮想列(ファイル名の一部であるyear列やshipmode列)は原理的に変わる事はあり得ません。
Arrow_Fdwがf_lineorder_1993_TRUCK.arrowのスキャンを終えたので、次にf_lineorder_1994_AIR.arrowを読み出そう、というタイミングで仮想列の内容が変化するだけですね。

Arrow_Fdw仮想列を使った「読み飛ばし」

では、このArrow_Fdw仮想列を使って問い合わせを最適化してみる事にします。

以下のクエリはStar Schema BenchmarkのQ1_1ですが、このベンチマークは日付の絞り込みにdate1テーブルとのJOINを行って、そのdate1テーブルのフィールドをd_year = 1993のように指定して絞り込みを行うため、Arrow_Fdw仮想列の絞り込みが効きません。

postgres=# explain
select sum(lo_extendedprice*lo_discount) as revenue
from f_lineorder,date1
where lo_orderdate = d_datekey
and d_year = 1993
and lo_discount between 1 and 3
and lo_quantity < 25;
                                                                         QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------------------------------
 Aggregate  (cost=178617.28..178617.29 rows=1 width=32)
   ->  Gather  (cost=178617.17..178617.28 rows=1 width=32)
         Workers Planned: 2
         ->  Parallel Custom Scan (GpuPreAgg) on f_lineorder  (cost=177617.17..177617.18 rows=1 width=32)
               GPU Projection: pgstrom.psum((lo_extendedprice * lo_discount))
               GPU Scan Quals: ((lo_discount >= '1'::numeric) AND (lo_discount <= '3'::numeric) AND (lo_quantity < '25'::numeric)) [plan: 65062080 -> 45182]
               GPU Join Quals [1]: (lo_orderdate = d_datekey) [plan: 45182 -> 6452]
               GPU Outer Hash [1]: lo_orderdate
               GPU Inner Hash [1]: d_datekey
               GPU Group Key:
               referenced: lo_orderdate, lo_quantity, lo_extendedprice, lo_discount
               file0: /opt/arrow/mydata/f_lineorder_1996_MAIL.arrow (read: 107.83MB, size: 427.16MB)
               file1: /opt/arrow/mydata/f_lineorder_1996_SHIP.arrow (read: 107.82MB, size: 427.13MB)
                   :                     :                       :
               file28: /opt/arrow/mydata/f_lineorder_1995_MAIL.arrow (read: 107.43MB, size: 425.58MB)
               file29: /opt/arrow/mydata/f_lineorder_1993_TRUCK.arrow (read: 107.51MB, size: 425.91MB)
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>
               ->  Parallel Seq Scan on date1  (cost=0.00..65.79 rows=215 width=4)
                     Filter: (d_year = 1993)
(19 rows)

しかし、lo_orderdateが1993年のものという条件で検索をするのであれば、以下のようにクエリを書き換えても良いはずです。
date1テーブルへのJOINを外し、代わりに仮想列であるyearを参照しyear = 1993という条件を付加しました。

postgres=# explain analyze
select sum(lo_extendedprice*lo_discount) as revenue
from f_lineorder
where year = 1993
and lo_discount between 1 and 3
and lo_quantity < 25;
                                                                                               QUERY PLAN
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 Aggregate  (cost=167002.79..167002.80 rows=1 width=32) (actual time=74.534..74.614 rows=1 loops=1)
   ->  Gather  (cost=167002.67..167002.78 rows=1 width=32) (actual time=74.516..74.600 rows=2 loops=1)
         Workers Planned: 2
         Workers Launched: 2
         ->  Parallel Custom Scan (GpuPreAgg) on f_lineorder  (cost=166002.67..166002.68 rows=1 width=32) (actual time=25.371..25.372 rows=1 loops=3)
               GPU Projection: pgstrom.psum((lo_extendedprice * lo_discount))
               GPU Scan Quals: ((year = 1993) AND (lo_discount <= '3'::numeric) AND (lo_quantity < '25'::numeric) AND (lo_discount >= '1'::numeric)) [plan: 65062080 -> 226, exec: 13010375 -> 1703647]
               GPU Group Key:
               referenced: lo_quantity, lo_extendedprice, lo_discount, year
               Stats-Hint: (year = 1993)  [loaded: 12, skipped: 48]
               file0: /opt/arrow/mydata/f_lineorder_1996_MAIL.arrow (read: 99.53MB, size: 427.16MB)
               file1: /opt/arrow/mydata/f_lineorder_1996_SHIP.arrow (read: 99.52MB, size: 427.13MB)
                   :                     :                       :
               file28: /opt/arrow/mydata/f_lineorder_1995_MAIL.arrow (read: 99.16MB, size: 425.58MB)
               file29: /opt/arrow/mydata/f_lineorder_1993_TRUCK.arrow (read: 99.24MB, size: 425.91MB)
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=76245, ntuples=13010375
 Planning Time: 0.718 ms
 Execution Time: 75.027 ms
(18 rows)

そうすると、EXPLAINの出力にはStats-Hint: (year = 1993) [loaded: 12, skipped: 48]という出力が表示されています。
これは、条件句にyear = 1993での絞り込みが含まれており、ファイル名に 1993 を含まないものは明らかに検索条件に合致しないために読み飛ばしているという事を意味します。
これらのArrowファイルはそれぞれ425MB前後とそれほど大きくはないですが、1ファイルあたり2個のRecord-Batchを含んでおり、1993年のファイル6個に対しては12個のRecord-Batchを読み出す必要があったという事を意味しています。

ここで利用できるのは等価演算子だけでなく、例えば不等号であっても読み飛ばしを行うことは可能です。

この機能は、特に非常に多くのファイルを一定の命名規則の下に管理している場合には非常に有効で、実際に2000個を越えるArrowファイルを仮想列の機能を用いて管理しているといった事例も存在します。

PG-Strom v6.0特集:GPU-Sortと一部のWindow関数に対応(その2)

PG-Strom v6.0をリリースしました。

  • GPU-Sortと一部のWindow関数対応
  • マルチGPUのPinned Inner Buffer対応
  • Arrow_Fdwの仮想列機能
  • GPUでの完全な集計値の生成

といった、いくつかの重要機能を含むリリースで、特にGPU-Sortによって新しいワークロードへの対応が広がったという事でバージョン6.x系列としました。
その主要機能について、何回かに分けて解説していきたいと思います。

前回の記事では、GPU-SortをGPU-Joinにアタッチし、LIMIT句のプッシュダウンによりGPUからCPUへ返却する行数を減らす処理について解説しました。
今回は、その続きとして集約関数を実行後にソートする処理、およびWindow関数のプッシュダウンについて説明します。

完全なAggregationをGPUで生成する

以下の実行計画を見てください。
4つのテーブル(lineorder、date1、part、supplier)をJOINし、d_year列、p_brand1列によるGROUP BYと集約関数AVG()を実行するシンプルな問い合わせです。

=# explain
     select avg(lo_revenue), d_year, p_brand1
       from lineorder, date1, part, supplier
      where lo_orderdate = d_datekey
        and lo_partkey = p_partkey
        and lo_suppkey = s_suppkey
        and p_category = 'MFGR#12'
        and s_region = 'AMERICA'
      group by d_year, p_brand1;
                                                                            QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------------------------------------
 HashAggregate  (cost=13806407.68..13806495.18 rows=7000 width=46) (actual time=43407.303..43410.313 rows=280 loops=1)
   Group Key: date1.d_year, part.p_brand1
   Batches: 1  Memory Usage: 273kB
   ->  Gather  (cost=13805618.73..13806355.18 rows=7000 width=46) (actual time=43405.470..43409.956 rows=560 loops=1)
         Workers Planned: 2
         Workers Launched: 2
         ->  Parallel Custom Scan (GpuPreAgg) on lineorder  (cost=13804618.73..13804655.18 rows=7000 width=46) (actual time=43396.506..43396.532 rows=187 loops=3)
               GPU Projection: pgstrom.pavg(lo_revenue), d_year, p_brand1
               GPU Join Quals [1]: (p_partkey = lo_partkey) [plan: 2500011000 -> 98584180, exec: 7520 -> 250]
               GPU Outer Hash [1]: lo_partkey
               GPU Inner Hash [1]: p_partkey
               GPU Join Quals [2]: (s_suppkey = lo_suppkey) [plan: 98584180 -> 19644550, exec: 250 -> 0]
               GPU Outer Hash [2]: lo_suppkey
               GPU Inner Hash [2]: s_suppkey
               GPU Join Quals [3]: (d_datekey = lo_orderdate) [plan: 19644550 -> 19644550, exec: 0 -> 0]
               GPU Outer Hash [3]: lo_orderdate
               GPU Inner Hash [3]: d_datekey
               GpuJoin buffer usage: 143.68MB
               GPU Group Key: d_year, p_brand1
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=7520
               ->  Parallel Custom Scan (GpuScan) on part  (cost=100.00..12682.86 rows=32861 width=14) (actual time=20.596..45.648 rows=26528 loops=3)
                     GPU Projection: p_brand1, p_partkey
                     GPU Scan Quals: (p_category = 'MFGR#12'::bpchar) [plan: 2000000 -> 32861, exec: 2000000 -> 79584]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=29258, ntuples=2000000
               ->  Parallel Custom Scan (GpuScan) on supplier  (cost=100.00..78960.47 rows=830255 width=6) (actual time=10.240..66.126 rows=667164 loops=3)
                     GPU Projection: s_suppkey
                     GPU Scan Quals: (s_region = 'AMERICA'::bpchar) [plan: 9999718 -> 830255, exec: 10000000 -> 2001491]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=168663, ntuples=10000000
               ->  Parallel Seq Scan on date1  (cost=0.00..62.04 rows=1504 width=8) (actual time=0.008..0.139 rows=852 loops=3)
 Planning Time: 0.755 ms
 Execution Time: 43410.797 ms
(31 rows)

この実行計画では、CPUで集約関数を実行するHashAggregateの下に、並列ワーカープロセスを制御するGatherノード、そしてGPUで部分的な集計処理を行うGpuPreAggがぶら下がっています。
GpuPreAggが返却する結果セットはGPU Projection行にある通り、d_year列、p_brand1列、および集計値のpgstrom.pavg(lo_revenue)です。
集計値は非NULLなlo_revenue列の件数とその総和を含むバイナリデータ(bytea型)で、平均値の定義より、各ワーカープロセスがこれらの値を返せば HashAggregateは平均値を計算する事ができます。

では、なぜこのような2段階の方式を取っているのでしょうか?
これには歴史的な経緯と技術的な理由がそれぞれあります。

歴史的な経緯

PG-Strom v3.x系列まではPostgreSQLバックエンドプロセスが個別にCUDAコンテキストを管理していました。
これはGPUの世界における「プロセス」のようなもので、あるCUDAコンテキストから獲得したGPUメモリを別のCUDAコンテキストから参照するにはひと手間必要でした。

そうすると、PostgreSQLバックグラウンドワーカーで作成したCUDAコンテキスト同士で結果をマージするよりは、GPUで部分的な集計を行い、ワーカーの終了後にCPUで並列処理結果を集計した方がシンプルな構造だったわけです。(これはPostgreSQLの並列クエリと同じ処理の分割方法です)

技術的な理由

もう一つは、CPU-Fallbackに関連するものです。(また出た😅)
典型的にはGPUでの処理中にTOAST化された可変長データを参照した場合など、エラーではないもののGPUでの処理を継続できない場合に、PG-Stromはその行をCPU側に書き戻して、CPUで再実行するCPU-Fallbackという機能を持っています。この可能性がある事で、集計値を計算する際にCPU側にFallbackされた未集計のデータが存在する事を排除できず、そうすると結局はCPU側で最終的な集計処理を行わざるを得ないよねという事になってしまいます。

そういうわけで、Pinned Inner Buffer機能と同じく、明示的にCPU-Fallbackがoffの場合に限って、GPU上で全ての集計処理を行ってしまうという割り切りが必要になったわけです。

完全なAggregationをGPU側で作成する。

以下の実行計画を見てください。(説明のため verbose モードにしています)

ssbm=# set pg_strom.cpu_fallback = off;
SET
ssbm=# explain (verbose, analyze)
select avg(lo_revenue), d_year, p_brand1
from lineorder, date1, part, supplier
where lo_orderdate = d_datekey
and lo_partkey = p_partkey
and lo_suppkey = s_suppkey
and p_category = 'MFGR#12'
and s_region = 'AMERICA'
  group by d_year, p_brand1;
                                                                               QUERY PLAN
-------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 Gather  (cost=13805618.73..13806460.18 rows=7000 width=46) (actual time=43234.075..43246.038 rows=280 loops=1)
   Output: (pgstrom.favg_numeric((pgstrom.pavg(lineorder.lo_revenue)))), date1.d_year, part.p_brand1
   Workers Planned: 2
   Workers Launched: 2
   ->  Result  (cost=13804618.73..13804760.18 rows=7000 width=46) (actual time=43224.472..43224.530 rows=93 loops=3)
         Output: pgstrom.favg_numeric((pgstrom.pavg(lineorder.lo_revenue))), date1.d_year, part.p_brand1
         Worker 0:  actual time=43218.832..43218.837 rows=0 loops=1
         Worker 1:  actual time=43228.887..43229.052 rows=280 loops=1
         ->  Parallel Custom Scan (GpuPreAgg) on public.lineorder  (cost=13804618.73..13804655.18 rows=7000 width=46) (actual time=43224.466..43224.482 rows=93 loops=3)
               Output: (pgstrom.pavg(lineorder.lo_revenue)), date1.d_year, part.p_brand1
               GPU Projection: pgstrom.pavg(lineorder.lo_revenue), date1.d_year, part.p_brand1
               GPU Join Quals [1]: (part.p_partkey = lineorder.lo_partkey) [plan: 2500011000 -> 98584180, exec: 7520 -> 250]
               GPU Outer Hash [1]: lineorder.lo_partkey
               GPU Inner Hash [1]: part.p_partkey
               GPU Join Quals [2]: (supplier.s_suppkey = lineorder.lo_suppkey) [plan: 98584180 -> 19644550, exec: 250 -> 0]
               GPU Outer Hash [2]: lineorder.lo_suppkey
               GPU Inner Hash [2]: supplier.s_suppkey
               GPU Join Quals [3]: (date1.d_datekey = lineorder.lo_orderdate) [plan: 19644550 -> 19644550, exec: 0 -> 0]
               GPU Outer Hash [3]: lineorder.lo_orderdate
               GPU Inner Hash [3]: date1.d_datekey
               GpuJoin buffer usage: 143.68MB
               GPU Group Key: date1.d_year, part.p_brand1
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=7520
               Worker 0:  actual time=43218.831..43218.836 rows=0 loops=1
               Worker 1:  actual time=43228.872..43228.910 rows=280 loops=1
               ->  Parallel Custom Scan (GpuScan) on public.part  (cost=100.00..12682.86 rows=32861 width=14) (actual time=4.688..21.335 rows=26528 loops=3)
                     Output: part.p_brand1, part.p_partkey
                     GPU Projection: part.p_brand1, part.p_partkey
                     GPU Scan Quals: (part.p_category = 'MFGR#12'::bpchar) [plan: 2000000 -> 32861, exec: 2000000 -> 79584]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=29258, ntuples=2000000
                     Worker 0:  actual time=0.729..0.732 rows=0 loops=1
                     Worker 1:  actual time=0.805..0.806 rows=0 loops=1
               ->  Parallel Custom Scan (GpuScan) on public.supplier  (cost=100.00..78960.47 rows=830255 width=6) (actual time=35.323..92.328 rows=667164 loops=3)
                     Output: supplier.s_suppkey
                     GPU Projection: supplier.s_suppkey
                     GPU Scan Quals: (supplier.s_region = 'AMERICA'::bpchar) [plan: 9999718 -> 830255, exec: 10000000 -> 2001491]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=168663, ntuples=10000000
                     Worker 0:  actual time=92.302..181.533 rows=1067881 loops=1
                     Worker 1:  actual time=13.335..95.117 rows=933610 loops=1
               ->  Parallel Seq Scan on public.date1  (cost=0.00..62.04 rows=1504 width=8) (actual time=0.007..0.111 rows=852 loops=3)
                     Output: date1.d_year, date1.d_datekey
                     Worker 0:  actual time=0.004..0.004 rows=0 loops=1
                     Worker 1:  actual time=0.004..0.004 rows=0 loops=1
 Planning Time: 0.688 ms
 Execution Time: 43246.432 ms
(45 rows)

注意して見ないと分かりませんが、Gather -> Result -> Parallel Custom Scan (GpuPreAgg) という流れでクエリが実行されており、最初の例に存在した Hash Aggregate つまりCPUでの集約処理が存在していません。
依然としてGpuPreAggが返却するのはpgstrom.pavg(lineorder.lo_revenue)というlo_revenue列の件数と総和のペアをパックしたバイナリ値なのですが、それを平均値に直してより上位のプランに出力するのはResultです。これは入力行を1行ずつ単純変換するだけの軽量なプランですので、実行性能への影響はほとんどありません。
そのため、CPUで最終盤の集計処理を行う方法と比べ、特にグループ数の大きな場合では性能改善効果が大きくなります。

そして、GpuPreAggの出力の中で、以下の行に注目してください。

               Worker 0:  actual time=43218.831..43218.836 rows=0 loops=1
               Worker 1:  actual time=43228.872..43228.910 rows=280 loops=1

このケースではGPUを2台使用し、そしてPostgreSQLバックエンドプロセスとワーカープロセスが2個の計3プロセスで実行していますが、今回はGpuPreAggのスキャンする lineorder 表の終端処理を行ったのが(たまたま)Worker-1で、Worker-1が最後にGPU0とGPU1の処理結果を統合し、それをCPU側に書き戻しているという事が分かります。

そして最も重要なのが、ここで集約演算を行った完全な結果セットを作り出す事ができたという事は、GPUメモリにデータが乗ったそのままの状態で、SortやWindow関数と言った「その後」の処理までをGPUで実行するための道が拓けたという事です。

Window関数のプッシュダウン

そして最後に、今回の目玉機能の一つ、Window関数のプッシュダウンについて説明します。
まず、従来の動作(CPU Fallback=on状態)で実行した以下のクエリと実行計画をご覧ください。

このクエリは、lineorderとcustomerテーブルをJOINし、c_region、c_nation、c_cityとlo_orderdate、sum(lo_revenue)を出力しますが、それぞれc_region、c_nation、c_cityごとにsum(lo_revenue)の合計値の多い日付を上位3件ずつ表示するというものです。分析系のクエリではそこそこ遭遇しそうなパターンの気がします。

実行計画は、GpuPreAggが部分集計結果を返し、GatherとHashAggregateがそれを集約。この時点で60,1500行が中間出力されています。
その後、Run Condition: (rank() OVER (?) < 4)を含むWindowAggによって、各パーティションの上位3件で計750件のデータが出力されています。

言い換えれば、それ以外の60万行は最終結果には寄与しないにも関わらず、GPUからCPUへ書き戻され、ご丁寧にそれをHashAggregateやSortで集計・整列の処理を行っているわけです。

ssbm=# explain analyze
   select * from (
       select c_region, c_nation, c_city, lo_orderdate, sum(lo_revenue) lo_rev,
              rank() over(partition by c_region, c_nation, c_city
                          order by sum(lo_revenue)) cnt
         from lineorder, customer
        where lo_custkey = c_custkey
         group by c_region, c_nation, c_city, lo_orderdate
   ) subqry
   where cnt < 4;
                                                                                      QUERY PLAN
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 WindowAgg  (cost=74232835.17..76112522.67 rows=75187500 width=84) (actual time=58176.445..58352.400 rows=750 loops=1)
   Run Condition: (rank() OVER (?) < 4)
   ->  Sort  (cost=74232835.17..74420803.92 rows=75187500 width=76) (actual time=58176.435..58256.419 rows=601500 loops=1)
         Sort Key: customer.c_region, customer.c_nation, customer.c_city, (pgstrom.sum_numeric((pgstrom.psum(lineorder.lo_revenue))))
         Sort Method: quicksort  Memory: 76268kB
         ->  HashAggregate  (cost=58207057.50..61055958.87 rows=75187500 width=76) (actual time=55033.033..55321.629 rows=601500 loops=1)
               Group Key: customer.c_region, customer.c_nation, customer.c_city, lineorder.lo_orderdate
               Planned Partitions: 8  Batches: 1  Memory Usage: 516113kB
               ->  Gather  (cost=40216881.72..48127233.28 rows=75187500 width=76) (actual time=53907.157..54301.344 rows=1203000 loops=1)
                     Workers Planned: 2
                     Workers Launched: 2
                     ->  Parallel Custom Scan (GpuPreAgg) on lineorder  (cost=40215881.72..40607483.28 rows=75187500 width=76) (actual time=53808.628..53903.530 rows=401000 loops=3)
                           GPU Projection: pgstrom.psum(lo_revenue), c_region, c_nation, c_city, lo_orderdate
                           GPU Join Quals [1]: (lo_custkey = c_custkey) [plan: 2500011000 -> 2500011000, exec: 3715510 -> 3413504]
                           GPU Outer Hash [1]: lo_custkey
                           GPU Inner Hash [1]: c_custkey
                           GpuJoin buffer usage: 3204.35MB
                           GPU Group Key: c_region, c_nation, c_city, lo_orderdate
                           Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=3715510
                           ->  Parallel Seq Scan on customer  (cost=0.00..644671.73 rows=12501573 width=46) (actual time=0.031..1472.440 rows=10000000 loops=3)
 Planning Time: 0.449 ms
 Execution Time: 58734.809 ms
(22 rows)

では次に、CPU-Fallbackを無効化した状態で同じクエリを実行してみます。

ssbm=# set pg_strom.cpu_fallback = off;
SET
ssbm=# explain analyze
   select * from (
       select c_region, c_nation, c_city, lo_orderdate, sum(lo_revenue) lo_rev,
              rank() over(partition by c_region, c_nation, c_city
                          order by sum(lo_revenue)) cnt
         from lineorder, customer
        where lo_custkey = c_custkey
         group by c_region, c_nation, c_city, lo_orderdate
   ) subqry
   where cnt < 4;
                                                                                QUERY PLAN
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 WindowAgg  (cost=40296711.63..40315461.63 rows=125000 width=84) (actual time=54055.549..54373.242 rows=750 loops=1)
   Run Condition: (rank() OVER (?) < 4)
   ->  Gather  (cost=40296711.63..40312649.13 rows=125000 width=76) (actual time=54055.534..54372.279 rows=750 loops=1)
         Workers Planned: 2
         Workers Launched: 2
         ->  Result  (cost=40295711.63..40299149.13 rows=125000 width=76) (actual time=53960.565..53960.756 rows=250 loops=3)
               ->  Parallel Custom Scan (GpuPreAgg) on lineorder  (cost=40295711.63..40297274.13 rows=125000 width=76) (actual time=53960.559..53960.635 rows=250 loops=3)
                     GPU Projection: pgstrom.psum(lo_revenue), c_region, c_nation, c_city, lo_orderdate
                     GPU Join Quals [1]: (lo_custkey = c_custkey) [plan: 2500011000 -> 2500011000, exec: 3749497 -> 3444736]
                     GPU Outer Hash [1]: lo_custkey
                     GPU Inner Hash [1]: c_custkey
                     GpuJoin buffer usage: 3204.35MB
                     GPU Group Key: c_region, c_nation, c_city, lo_orderdate
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=3749497
                     GPU-Sort keys: c_region, c_nation, c_city, pgstrom.fsum_numeric((pgstrom.psum(lo_revenue)))
                     Window-Rank Filter: rank() over(PARTITION BY c_region, c_nation, c_city ORDER BY pgstrom.fsum_numeric((pgstrom.psum(lo_revenue)))) < 4
                     ->  Parallel Seq Scan on customer  (cost=0.00..644671.73 rows=12501573 width=46) (actual time=0.032..1463.581 rows=10000000 loops=3)
 Planning Time: 0.401 ms
 Execution Time: 54707.476 ms
(19 rows)

なかなかキモイ感じでシンプル化されているのがお分かりでしょうか?

Window関数を含むクエリの場合、下位プランに『~~の順にソートされていると嬉しいよ』とヒントが与えられます。
下位プランの作成・評価の際には、それを見た上で、ドンピシャのB-treeインデックスが存在すればそれを使いますし、Sortのプランを挟むといったアレンジも行われます。

このGpuPreAggでは、まずCPU-Fallback=offである事から完全な集約関数の結果を保持している事を前提として、rank()関数のPARTITION BY句でc_region, c_nation, c_cityが指定されている事から、GPU-Sortをアタッチした方がコストの低いプランを作れます。以下のようになっているわけですね。

GPU-Sort keys: c_region, c_nation, c_city, pgstrom.fsum_numeric((pgstrom.psum(lo_revenue)))

さらに、Window関数の使い方がrank() < CONSTである事から、これは「上位xx件」という形で行数を削除するタイプのクエリである事を判定します。
そうすると、rank()以外のよりバリエーションのあるWindow関数の処理はCPU側に任せるとしても、少なくともフィルタリングの対象となる行をわざわざCPUに戻して負荷を上げる事は避けられるというわけです。

Window-Rank Filter: rank() over(PARTITION BY c_region, c_nation, c_city ORDER BY pgstrom.fsum_numeric((pgstrom.psum(lo_revenue)))) < 4

現在のところ、以下のパターンのWindow関数のプッシュダウンに対応しています。

  • rank() < または <= CONST
  • dense_rank() < または <= CONST
  • row_number() < または <= CONST

これらの機能により、分析系クエリではそこそこ使用頻度の高くチューニングの方法がほとんどないWindow関数に関しても、やりようによっては高速化のアプローチが一つ手に入ったという事になるのではないでしょうか。

PG-Strom v6.0特集:GPU-Sortと一部のWindow関数に対応(その1)

PG-Strom v6.0をリリースしました。

  • GPU-Sortと一部のWindow関数対応
  • マルチGPUのPinned Inner Buffer対応
  • Arrow_Fdwの仮想列機能
  • GPUでの完全な集計値の生成

といった、いくつかの重要機能を含むリリースで、特にGPU-Sortによって新しいワークロードへの対応が広がったという事でバージョン6.x系列としました。
その主要機能について、何回かに分けて解説していきたいと思います。

GPU-Sortの歴史的経緯

長らく、PG-StromがオフロードするワークロードはSCAN、JOIN、GROUP-BYの三種類でした。
これらにGPUバイスで対応しているSQL関数や演算子の組み合わせをGPUで並列実行し、またGPU-Direct SQLによる高速なI/Oサブシステムによって、大量のデータをスキャンし、集計した結果を出力するというのがPG-Stromの使いどころでした。

これはそれほど根拠レスな話ではありません。以前に分析系SQLの性能評価に使われるTPC-DSのクエリ103本の実行時間を調査して、SQLのワークロードのうちにどのような処理が実際に時間を費やしているのか調べた事があるのですが、ざっくり34%がSCAN、37%がJOIN、23%が集計処理という事で、この3つで総処理時間の95%を占めています。

しかし、行数が多い場合には特異的にSort処理の時間がかかったり、またSortを前処理として要求するWindow関数の処理に時間を要したりしているクエリもあり、より網羅的なワークロードへの対応という事では選択肢の一つでした。

実はかつて、PG-StromでもGPU-Sortの開発を試みた事がありました。
PG-Stromのリポジトリには、『もはや使わなくなったけれども、もしかしたら同じロジックを将来使うかもしれないし、その時の参考になるかもしれない』コードをモスポール保存しておくためのdeadcodeディレクトリというものがあるのですが、そこにopencl_gpusort.cというコードが存在するほどです。
github.com

CUDAではなくOpenCLを実装のベースとしていた頃なので、おそらく2015年前後。今からちょうど10年前だと思います。
この頃、開発に使用していたGPUはGTX980(メモリ4GB)やTESLA K20(メモリ5GB)といったモデルで、GPUに搭載可能なデータ量というのは今とは比較にならないほど小さなものでした。物理的な限界もさておき、当時はManaged Memoryを使用した時でもGPU Memory Oversubscriptionができなかったため、バッファサイズの推定(予想)はかなり保守的に行わざるを得ず、3GB程度のバッファを用意してスキャンを開始しても、実際にテーブルから読み出して条件句に合致する行だけをバッファに乗せたら200MB程度しか使っていなかった、という事もザラでした。

そうすると、GPUでソート処理を行ったとしても大した規模のソート処理を行うことはできず、最終的にCPUでマージソートを実行せねばならないため、それなら最初からCPUのメモリに乗せてクイックソートでもした方が速いのでは・・・みたいな事になる事もザラでした。実際、当時のGPU-Sortは良くてCPUと同じくらい。大抵は遅くなるというシロモノでした。

しかし、最近のGPUは急速に搭載メモリを増やしてきています。2017年にTESLA V100の搭載メモリが16GBだったのが、2020年のNVIDIA A100では40GBに、2023年のNVIDIA H100では80GBを搭載し、先日のGTC2025で発表されたRTX Pro 6000 Blackwellでは96GBのメモリが搭載されるようになっています。
これだけのデータをGPU上で一度に処理できるようになってくると話が変わってきて、GPU-JoinやGPU-PreAggの後処理として『ソート済み』の結果をCPUに返し、CPUでのソート処理を省略してやる事で、実用的な高速化の効果を実感できるようになってきたのです。

GPUメモリ上での処理フロー

Pinned Inner Bufferの機能とも関係するところですが、PG-StromがSCANやJOINのワークロードを処理する場合と、GROUP-BYのワークロードを処理する場合では、少しデータの動き方が異なります。

全てのワークロードに共通して、テーブルをスキャンする時には約64MB単位のブロックに区切って(Apache Arrowの場合はRecord-Batchを単位として)これをGPUメモリに読み出します。GPU上では命令コードに従ってWHERE句やJOIN結合条件の評価、あるいは集約関数の集計処理といった処理を行い、最後に処理結果を結果バッファに書き込みます。

基本的に、SCAN/JOINは読み出した64MBのブロックごとに結果を書き戻し(one-by-one execution)て、使用したバッファは次のブロックのために解放し再利用します。そのため、たとえテーブルの大きさが1TBあろうとも、少ないGPUメモリでこれを処理できます。
一方で、GROUP BYの場合は集計値をGPUメモリに置いておき、次の入力ブロック、次の入力ブロック、、、、と次々に集計値を保持する結果バッファを更新して、最後に結果を書き戻す事になります。

ソート処理を実行するには、対象となるデータが全てGPUメモリ上に乗っていることが必要です。
そのため、SCANやJOINにGPU-Sortが付加する場合は、これまでのGROUP-BYのインフラを利用して、処理結果をGPUに留置するように手が加えられています。つまり、SCANやJOINであっても Dam-execution 方式となるため、結果セットがGPUメモリの範囲内に収まるかどうかは注意が必要です。

もう一点、これはPinned Inner Bufferとも共通の制限事項ですが、CPU-Fallback機構と同時に利用する事はできません。
例えば非常に長い可変長データがTOAST化されていた場合など、GPUでその文字列やジオメトリデータを処理したくとも、GPUにロードされたデータの中にそれが含まれていなければどうしようもありません。そういった場合には、CPUで再実行するために対象行をCPU側へ書き戻し、CPUで再実行する仕組みがPG-Stromには備わっています。これをCPU-Fallback機構と呼んでいるのですが、この仕組みによって行がCPUに書き戻された場合、GPUメモリ上にはソート処理に必要な「完全な結果セット」が存在しない事になります。

そのため、GPU-Sortの有効化には、ワークロードが結果セットの並び替えを要求している事と同時に、CPU-Fallbackが無効化されている事が必要です。

例えば、以下のクエリ(ORDER BYを含む)は GpuJoin の処理結果に部分Sort(これはGather配下のワーカープロセスでも実行されるため)を実行し、各ワーカープロセスが生成した部分Sort済みの結果セットを Gather Merge で統合しているのが分かります。そして、最後にLimitを実行し上位15件のみを取り出します。
全てのソート処理はCPUで実行されているため、GPUでの処理結果(GpuJoin)は特に並び替えられていません。

ssbm=# explain analyze
        select d_year, p_brand1
          from lineorder, date1, part, supplier
         where lo_orderdate = d_datekey
           and lo_partkey = p_partkey
           and lo_suppkey = s_suppkey
           and p_brand1 in ('MFGR#2221','MFGR#2228')
           and s_region in ('ASIA','AMERICA')
         order by d_year, lo_discount
         limit 15;
                                                                                QUERY PLAN
---------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 Limit  (cost=12569927.20..12569928.96 rows=15 width=18) (actual time=44083.849..44103.448 rows=15 loops=1)
   ->  Gather Merge  (cost=12569927.20..13031121.74 rows=3952820 width=18) (actual time=44083.848..44103.444 rows=15 loops=1)
         Workers Planned: 2
         Workers Launched: 2
         ->  Sort  (cost=12568927.18..12573868.21 rows=1976410 width=18) (actual time=44075.562..44075.567 rows=15 loops=3)
               Sort Key: date1.d_year, lineorder.lo_discount
               Sort Method: top-N heapsort  Memory: 26kB
               Worker 0:  Sort Method: top-N heapsort  Memory: 27kB
               Worker 1:  Sort Method: top-N heapsort  Memory: 26kB
               ->  Parallel Custom Scan (GpuJoin) on lineorder  (cost=143664.84..12520437.04 rows=1976410 width=18) (actual time=1070.804..43861.693 rows=1453231 loops=3)
                     GPU Projection: d_year, p_brand1, lo_discount
                     GPU Join Quals [1]: (p_partkey = lo_partkey) [plan: 2500011000 -> 4976272, exec: 5999989709 -> 12086435]
                     GPU Outer Hash [1]: lo_partkey
                     GPU Inner Hash [1]: p_partkey
                     GPU Join Quals [2]: (s_suppkey = lo_suppkey) [plan: 4976272 -> 1976410, exec: 12086435 -> 4359693]
                     GPU Outer Hash [2]: lo_suppkey
                     GPU Inner Hash [2]: s_suppkey
                     GPU Join Quals [3]: (d_datekey = lo_orderdate) [plan: 1976410 -> 1976410, exec: 4359693 -> 4359693]
                     GPU Outer Hash [3]: lo_orderdate
                     GPU Inner Hash [3]: d_datekey
                     GpuJoin buffer usage: 275.22MB
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=5999989709
                     ->  Parallel Custom Scan (GpuScan) on part  (cost=100.00..12370.84 rows=1659 width=14) (actual time=22.416..50.698 rows=1337 loops=3)
                           GPU Projection: p_brand1, p_partkey
                           GPU Scan Quals: (p_brand1 = ANY ('{MFGR#2221,MFGR#2228}'::bpchar[])) [plan: 2000000 -> 1659, exec: 2000000 -> 4010]
                           Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=29258, ntuples=2000000
                     ->  Parallel Custom Scan (GpuScan) on supplier  (cost=100.00..87206.07 rows=1654815 width=6) (actual time=8.979..120.408 rows=1333704 loops=3)
                           GPU Projection: s_suppkey
                           GPU Scan Quals: (s_region = ANY ('{ASIA,AMERICA}'::bpchar[])) [plan: 9999718 -> 1654815, exec: 10000000 -> 4001111]
                           Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=168663, ntuples=10000000
                     ->  Parallel Seq Scan on date1  (cost=0.00..62.04 rows=1504 width=8) (actual time=0.008..0.112 rows=852 loops=3)
 Planning Time: 0.970 ms
 Execution Time: 44103.741 ms
(33 rows)

ここで、CPU-Fallbackを無効化します。
実行計画がシンプルになったのがお分かりでしょうか?

ssbm=# set pg_strom.cpu_fallback = off;
SET
ssbm=# explain analyze
        select d_year, p_brand1
          from lineorder, date1, part, supplier
         where lo_orderdate = d_datekey
           and lo_partkey = p_partkey
           and lo_suppkey = s_suppkey
           and p_brand1 in ('MFGR#2221','MFGR#2228')
           and s_region in ('ASIA','AMERICA')
         order by d_year, lo_discount
         limit 15;
                                                                       QUERY PLAN
--------------------------------------------------------------------------------------------------------------------------------------------------------
 Gather  (cost=12514590.37..12514591.95 rows=15 width=18) (actual time=44147.437..44177.058 rows=15 loops=1)
   Workers Planned: 2
   Workers Launched: 2
   ->  Parallel Custom Scan (GpuJoin) on lineorder  (cost=12513590.37..12513590.45 rows=15 width=18) (actual time=43961.197..43961.204 rows=5 loops=3)
         GPU Projection: lo_discount, d_year, p_brand1
         GPU Join Quals [1]: (p_partkey = lo_partkey) [plan: 2500011000 -> 4976272, exec: 6000184056 -> 12086748]
         GPU Outer Hash [1]: lo_partkey
         GPU Inner Hash [1]: p_partkey
         GPU Join Quals [2]: (s_suppkey = lo_suppkey) [plan: 4976272 -> 1976410, exec: 12086748 -> 4359793]
         GPU Outer Hash [2]: lo_suppkey
         GPU Inner Hash [2]: s_suppkey
         GPU Join Quals [3]: (d_datekey = lo_orderdate) [plan: 1976410 -> 1976410, exec: 4359793 -> 4359793]
         GPU Outer Hash [3]: lo_orderdate
         GPU Inner Hash [3]: d_datekey
         GpuJoin buffer usage: 275.22MB
         Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=114826068, ntuples=6000184056
         GPU-Sort keys: d_year, lo_discount
         GPU-Sort Limit: 15
         ->  Parallel Custom Scan (GpuScan) on part  (cost=100.00..12370.84 rows=1659 width=14) (actual time=26.481..35.863 rows=1337 loops=3)
               GPU Projection: p_brand1, p_partkey
               GPU Scan Quals: (p_brand1 = ANY ('{MFGR#2221,MFGR#2228}'::bpchar[])) [plan: 2000000 -> 1659, exec: 2000000 -> 4010]
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=29258, ntuples=2000000
         ->  Parallel Custom Scan (GpuScan) on supplier  (cost=100.00..87206.07 rows=1654815 width=6) (actual time=12.544..122.628 rows=1333704 loops=3)
               GPU Projection: s_suppkey
               GPU Scan Quals: (s_region = ANY ('{ASIA,AMERICA}'::bpchar[])) [plan: 9999718 -> 1654815, exec: 10000000 -> 4001111]
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=168663, ntuples=10000000
         ->  Parallel Seq Scan on date1  (cost=0.00..62.04 rows=1504 width=8) (actual time=0.008..0.117 rows=852 loops=3)
 Planning Time: 0.812 ms
 Execution Time: 44187.259 ms
(29 rows)

GpuJoinのオプション列を見てみるとGPU-Sort keys: d_year, lo_discountと、GPU-Sort Limit: 15という表示が出ている事が分かるでしょうか?
この表示は、GpuJoinの処理結果をキー値d_yearとlo_discountに基づいてソートする事を示し、また、次の行はソートした結果のうち上位15件だけをGPUからCPUへ書き戻す事を示しています。
このクエリは元々430万件程度しか返しませんので、実行結果のGPU->CPU転送負荷というのはそれほど大したものではないのですが、件数が多くなってくるとプロセス間通信でそれをコピーするというのはそこそこ時間を要する処理となりますので、早い段階で不要な行を削るというのは合理的です。

この後、GPU-PreAggの改良とWindow関数のプッシュダウンについて説明したいと思いますが、長くなりましたのでエントリを分けようと思います。

PG-Strom v6.0特集:Large Tables JOIN(マルチGPU編)

PG-Strom v6.0をリリースしました。

  • GPU-Sortと一部のWindow関数対応
  • マルチGPUのPinned Inner Buffer対応
  • Arrow_Fdwの仮想列機能
  • GPUでの完全な集計値の生成

といった、いくつかの重要機能を含むリリースで、特にGPU-Sortによって新しいワークロードへの対応が広がったという事でバージョン6.x系列としました。
その主要機能について、何回かに分けて解説していきたいと思います。

マルチGPU環境でのPinned Inner Bufferの動きは?

前回、Large Tables JOIN を高速化する Pinned Inner Buffer について解説しました。

kaigai.hatenablog.com

これは、GPU-Scanの結果をCPU側へ戻さず、そのままGPUメモリに留置する事でGPU-Joinで使用するInner BufferをCPUで構築する処理を省略するという技術ですが、では、マルチGPUの環境ではどうなるのでしょうか?
なかなか個人では難しいですが、AWSやAzure、OCIなどの大手クラウドでは最新GPUを8台搭載したモンスターマシンがお手頃価格*1で利用できますし、マルチGPUであればGPU-JoinのInner-Bufferもより大きなサイズを利用できるのではないかというのは当然の期待でしょう。

先ず、大前提としてPG-StromはGPU-JoinのInner-BufferをCUDA Managed Memoryを用いて割り当てています。
この Managed Memory というのは10年以上前のKepler世代(CC3.0)から存在する機能で、CPUとGPUで同じ論理アドレス空間を使用し、その論理アドレス空間を読み書きした際に、もしも対応する物理ページが存在しなくとも、Page Faultを使ってGPUからCPUへ、CPUからGPUへとページ単位でコピーを行い、あたかも1個のメモリ領域をCPUとGPUで共有するための機能です。

Pascal世代(CC6.0)以降では、Manager Memoryをアロケートした時点では論理アドレス空間のみを割り当て、そのページを読み書きするたびにPage Faultを発生させて物理ページを割り当てるため、行列演算などと異なり、予め処理結果を保存するためのバッファサイズを読みにくいSQLのワークロードには大変ありがたい。
しかも、GPUのデバイスメモリが溢れた場合には、あたかもOSのスワップのように、GPUメモリの内容をCPUメモリに退避してくれるため、GPUメモリサイズ以上のバッファを消費しても処理を継続できます(GPU Memory Oversubscription)。

kaigai.hatenablog.com
実際、PG-Stromでもこの機能を便利に多用しており、これが『PG-StromはPascal世代以降でしか動作しない』と2017年の時点で旧世代のGPUを全てぶった切った最大の理由であります。

では、性能の方はどうでしょうか?
GpuJoinのPinned Inner Bufferを有効にして、前回のテストで用いたクエリの条件句o_orderdate > '1997-05-01'を調整してやれば、orders表のうちGpuJoin Inner Bufferに乗るデータサイズを調整する事ができます。

select l_shipmode, o_orderpriority, sum(l_extendedprice)
  from lineitem, orders
 where l_orderkey = o_orderkey
   and o_orderdate > '1997-05-01'
 group by l_shipmode, o_orderpriority;

そうして、Inner Bufferのサイズを徐々に大きくしていったものが以下のグラフです。横軸がInner Bufferのサイズ、縦軸がクエリの応答時間です。
すごいですね(棒)。このGPUのメモリサイズは40GBですが、GPU-Joinを実行するための最低限のワーキングとして4.2GBを要した上で、Pinned Inner Bufferのサイズが39GBを越えたあたりで急速にクエリの実行時間が伸びています。

合計すると44GBほどのメモリ消費で「壁」を越えられないように見えるかもしれませんが、実際にはInner Bufferの中でもGPU-Joinの処理中には滅多に参照されない行インデックスの領域が3.6GB(8byte x 4.7億行)ほどありますから、実態としてはGPUメモリのサイズを越えるようなメモリオーバーサブスクリプションは(特にHash-Joinのようなランダムアクセス性の強いものは)厳禁という事なのでしょう。

マルチGPUをどう活用するか

先ほどの実験より、Managed MemoryでGPUメモリ以上のバッファサイズを確保できるとは言えども、論理アドレスの世界からは見えなくしている物理的な制約からはどうしても逃れられない事が分かりました。つまり、単純に2GPUが使えるからと言って、倍のサイズのInner Bufferを確保しても性能が出るとは思えないのです。では、どうしたものでしょうか?

PG-Strom v6.0では、GPUメモリよりも大きな(実際にはpg_strom.pinned_inner_buffer_partition_sizeの設定値; GPUメモリの80~90%)Inner Bufferを複数のGPUに分割するよう試みます。
ここで重要なのは、GPU-Joinの実行中にできる限りPage Faultを避けるような分割を行わないと、GPU-DRAMに比べてもはるかに低速なPCI-Eバス上のデータ移動によってかなり破滅的な性能の低下が起こってしまうという点に留意する事です。つまり、JOIN処理で結合先の行を探す時には、そのGPUに存在しない行をハナから探しに行ってはいけないのです。(それでPage Faultを起こせばミイラ取りが・・・)

Inner Bufferの分割には、Hash-Joinで用いるハッシュ値を用います。
一般的にHash-Joinでinner_table.X = outer_table.Yという条件の結合を処理する際、outer_table.Yハッシュ値を計算し、予めハッシュ表にロードされたinner_tableのうち、inner_table.Xハッシュ値と一致するものとだけ値が同一かどうか検査します。
つまり、複数のGPUに分割してハッシュ表を格納する場合であっても、一定のルールに基づいて「このGPUにこのハッシュ値のINNER側の行は存在しない」という事を断言できるようにすれば、余計なPage Faultを引き起こしてしまう事はありません。

PG-Stromは、Hash-Joinで使用するハッシュ値パーティション数で割った剰余を使って場合分けをします。
前提として、OUTER側のテーブルから読み出した行は、ハッシュ値を計算したとしてもバラバラの値が混在しています。これはディスクから読み出したばかりのデータを、ディスク上に記録されている順序のまま読み出すのですから当然です。
しかし、Inner Hash Tableの方は予めJOIN結合キーが分かっており、それを元にGPUメモリ上で再編する事も可能です。

図のようなケースでは、Inner Bufferを3つのパーティションに分割し、各パーティションがそれぞれ3台のGPUにロードされています。
GPU0にはハッシュ値パーティション数(= 3)で割った剰余が0の行が、GPU1には剰余が1の、GPU2には剰余が2の行がそれぞれロードされています。

そして、ストレージから読み出したデータ(通常は約64MBのデータブロック)がGPU2にロードされると、まず、GPU2ではハッシュ値の剰余が2である行だけを、GPU1では剰余が1の、GPU0では剰余が0の行だけを処理します。
OUTER側の64MBブロックは、最初にGPU2にロードされると、次はGPU-to-GPUコピーによって隣のGPUへ次々と転送されていきます。しかし、GPU同士の通信は共にPCI-E x16レーンの広帯域を持っているためかなりのスピードが期待できますし、サーバーによってはSSD-to-GPUに用いるPCI-Eの経路ではなく、GPU同士を結合するNVLinkの通信を利用できるかもしれません。(こちらは未検証)

マルチGPUでのPinned Inner Bufferの動き

先ほどのTPC-Hデータセットを用いたクエリを利用して、マルチGPUでのPinned Inner Bufferの動きを観察してみます。

tpch=# set pg_strom.cpu_fallback = off;
SET
tpch=# set pg_strom.pinned_inner_buffer_threshold = '2GB';
SET
tpch=# explain analyze select l_shipmode, o_orderpriority, sum(l_extendedprice)
         from lineitem, orders
        where l_orderkey = o_orderkey
          and o_orderdate > '1996-01-01'
        group by l_shipmode, o_orderpriority;
NOTICE:  pinned inner-buffer partitions (depth=1, divisor=2)
NOTICE:  partition-0 (GPUs: 00000001)
NOTICE:  partition-1 (GPUs: 00000002)
                                                                           QUERY PLAN

----------------------------------------------------------------------------------------------------------------------------------------------------------------
 Gather  (cost=25490254.67..25490258.88 rows=35 width=59) (actual time=78253.459..78253.560 rows=35 loops=1)
   Workers Planned: 2
   Workers Launched: 2
   ->  Result  (cost=25489254.67..25489255.38 rows=35 width=59) (actual time=78241.277..78241.285 rows=12 loops=3)
         ->  Parallel Custom Scan (GpuPreAgg) on lineitem  (cost=25489254.67..25489254.85 rows=35 width=59) (actual time=78241.275..78241.280 rows=12 loops=3)
               GPU Projection: pgstrom.psum(l_extendedprice), l_shipmode, o_orderpriority
               GPU Join Quals [1]: (l_orderkey = o_orderkey) [plan: 2500102000 -> 979607700, exec: 548509 -> 94208]
               GPU Outer Hash [1]: l_orderkey
               GPU Inner Hash [1]: o_orderkey
               GpuJoin buffer usage: 4096B
               GPU Group Key: l_shipmode, o_orderpriority
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=115517748, ntuples=548509
               ->  Parallel Custom Scan (GpuScan) on orders  (cost=100.00..5285854.67 rows=244887632 width=24) (actual time=9926.267..9926.269 rows=-1 loops=3)
                     GPU Projection: o_orderpriority, o_orderkey
                     GPU Pinned Buffer: nitems: 588525682, usage: 48.23GB, total: 61.50GB, num-partitions: 2
                     GPU Scan Quals: (o_orderdate > '1996-01-01'::date) [plan: 1499973000 -> 244887600, exec: 1500000000 -> 588525682]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>; direct=26843246, ntuples=1500000000
 Planning Time: 0.331 ms
 Execution Time: 78253.739 ms
(19 rows)

GPU Pinned Buffer:行を見てみると、Pinned Bufferのサイズは61.50GBです。NVIDIA A100 [40GB; PCI-E] には物理的に載らないサイズですが、num-partitions: 2と2つに分割されている事が分かります。

実行中のI/Oの帯域と、GPU2台の使用率を見てみる事にします。

すると、ordersテーブルをスキャンしたと思しき処理の後、12秒程度、I/Oが休止している時間帯がある事が分かります。
これがHash値を元にPinned Inner Bufferの内容をGPU0とGPU1に再編している処理のタイミングで、決して無視できるレベルではないのですが、5.8億行、61.5GBのハッシュ表を移し替えるための処理としてはそこそこの所要時間と言えるでしょう。

次にGPU使用率を見てみます。order表のスキャン中はGPU0、GPU1ともに40%程度のGPU使用率に留まっており、ここはGPU負荷のそれほど大きくない処理である事が分かります。次いで、GPU0が100%に、そしてGPU1が100%となる時間帯がそれぞれ5~6秒ほどあり、ここでパーティションの再編を行っているものと思われます。その後、GPU0とGPU1が共に使用率80~90%程度で、ややI/Oの帯域が低下している時間帯が続きます。

ここはGPU-Joinの実行中で、GPU0での処理を終わった後にOUTER側テーブル(lineitem)から読み出したデータをGPU0->GPU1へ転送するために消費するPCI-Eバスの帯域、またはその逆にGPU1->GPU0へ転送するために消費するPCI-Eバスの帯域によって、ストレージから読み出しを行っている分の帯域が抑えられているものと思われます。

Inner Bufferのサイズ差による影響

最後に、Inner Bufferのサイズを色々と調整した場合の振る舞いを比較してみる事にします。

前述のように、orders表の検索条件をスライドする事でINNER側のバッファサイズを調整して試す事ができます。

検索条件 件数(百万) GPU-Join
Buffer-size
GPU-Join
実行時間
GPU-Join
Buffer-size
分割数 GPU-Join
実行時間
PostgreSQL
Hash-Size
PostgreSQL
実行時間
1998-01-01 132.80 10.88GB 65.28 14.00GB 1 53.07 7.44 262.94
1997-09-01 208.86 17.12GB 76.94 21.96GB 1 54.13 11.41 319.65
1997-05-01 285.54 23.40GB 89.85 29.95GB 1 56.28 15.42 377.57
1997-01-01 360.36 29.53GB 102.78 37.75GB 2 71.82 19.33 434.74
1996-09-01 436.41 35.77GB 116.22 45.67GB 2 73.58 23.31 492.16
1996-05-01 513.09 OOM ---- 53.65GB 2 77.40 27.31 551.69
1996-01-01 588.53 OOM ---- 61.50GB 2 79.84 31.25 608.58
1995-09-01 664.58 OOM ---- 69.44GB 3 189.41 35.23 670.02
1995-05-01 741.26 OOM ---- 77.42GB 3 189.81 39.24 727.64
1995-01-01 816.07 OOM ---- 85.21GB 3 231.76 43.14 784.58

従来のGPU-Joinの場合、GPUメモリサイズを越えるバッファの取扱いが十分でない事もあり、検索条件をo_orderdate > '1996-09-01'より緩くした場合のデータは測定できませんでした。
一方でPinned Inner Bufferを使用した場合、適切にバッファを分割する事で、検索条件を緩くしてもGPU-Joinを実行できるようになっています。
この場合、85GBのバッファと882GBのOUTER側テーブルとのJOINとなりますので、かなりの規模のテーブル同士でJOINを実行している事になります。

また、PostgreSQLの場合はCPU側でバッファを作成しており、実行時間自体はGPU-Joinよりかなり長いのですが、バッファサイズが小さい事に気が付かれたかもしれません。これは、歴史的な経緯などから、GPU-Joinに用いるハッシュ表のデータ形式が必ずしも十分にコンパクトでない(= 無駄がある)という事でもあります。

これは今後の課題と考えており、例えば1995-01-01のケースでCPU並みにバッファを削減できれば、43.14GBとなって余裕で2GPUに分割可能であり、その場合は80秒弱で当該のクエリを実行できるであろうことが期待できます。

先日のGTC2025では、Blackwell世代の新GPUであるNVIDIA RTX Pro 6000 Blackwellが発表されましたが、こちらはGPUメモリが96GBも搭載されており*2、こういった大容量RAMを搭載したGPUが普及するにつれてより使い勝手も良くなっていく事でしょう。

*1:諸説あります

*2:まぁ、生成AI向けがターゲットなんでしょうが

PG-Strom v6.0特集:Large Tables JOIN(シングルGPU編)

PG-Strom v6.0をリリースしました。

  • GPU-Sortと一部のWindow関数対応
  • マルチGPUのPinned Inner Buffer対応
  • Arrow_Fdwの仮想列機能
  • GPUでの完全な集計値の生成

といった、いくつかの重要機能を含むリリースで、特にGPU-Sortによって新しいワークロードへの対応が広がったという事でバージョン6.x系列としました。
その主要機能について、何回かに分けて解説していきたいと思います。

GPU-Joinの立ち上がり、遅くね?

以下の実行計画を見てみてください。

ssbm=# explain
select c_city, s_city, d_year, sum(lo_revenue) as revenue
from customer, lineorder, supplier, date1
where lo_custkey = c_custkey
and lo_suppkey = s_suppkey
and lo_orderdate = d_datekey
and c_nation = 'UNITED STATES'
and s_nation = 'UNITED STATES'
and d_year >= 1992 and d_year <= 1997
  group by c_city, s_city, d_year;
                                                   QUERY PLAN
-----------------------------------------------------------------------------------------------------------------
 HashAggregate  (cost=13553319.83..13558788.58 rows=437500 width=58)
   Group Key: customer.c_city, supplier.s_city, date1.d_year
   ->  Gather  (cost=13502916.18..13548944.83 rows=437500 width=58)
         Workers Planned: 2
         ->  Parallel Custom Scan (GpuPreAgg) on lineorder  (cost=13501916.18..13504194.83 rows=437500 width=58)
               GPU Projection: pgstrom.psum(lo_revenue), c_city, s_city, d_year
               GPU Join Quals [1]: (s_suppkey = lo_suppkey) [plan: 2500011000 -> 97750430]
               GPU Outer Hash [1]: lo_suppkey
               GPU Inner Hash [1]: s_suppkey
               GPU Join Quals [2]: (c_custkey = lo_custkey) [plan: 97750430 -> 3900243]
               GPU Outer Hash [2]: lo_custkey
               GPU Inner Hash [2]: c_custkey
               GPU Join Quals [3]: (d_datekey = lo_orderdate) [plan: 3900243 -> 3344810]
               GPU Outer Hash [3]: lo_orderdate
               GPU Inner Hash [3]: d_datekey
               GPU Group Key: c_city, s_city, d_year
               Scan-Engine: GPU-Direct with 2 GPUs <0,1>
               ->  Parallel Custom Scan (GpuScan) on supplier  (cost=100.00..72287.05 rows=162912 width=17)
                     GPU Projection: s_city, s_suppkey
                     GPU Scan Quals: (s_nation = 'UNITED STATES'::bpchar) [plan: 9999718 -> 162912]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>
               ->  Parallel Custom Scan (GpuScan) on customer  (cost=100.00..60041.62 rows=498813 width=17)
                     GPU Projection: c_city, c_custkey
                     GPU Scan Quals: (c_nation = 'UNITED STATES'::bpchar) [plan: 30003780 -> 498813]
                     Scan-Engine: GPU-Direct with 2 GPUs <0,1>
               ->  Parallel Seq Scan on date1  (cost=0.00..69.55 rows=1289 width=8)
                     Filter: ((d_year >= 1992) AND (d_year <= 1997))
(27 rows)

これはStar Schema BenchmarkのQ3_2ですが、実行計画のほぼ全体がGpuPreAgg内のJOINで完結しています。
最もサイズの大きな lineorder テーブルのスキャンを軸に、予めsupplier、customer、date1の各テーブルを読み出して構築したハッシュ表を用いてGPU上で並列のHash-Joinを実行する事が分かります。

このGpuPreAggは内部的にどのような動作をするかというと、CustomScan APIを通じてこのノードに処理が移ると

  1. ExecProcNodeを呼び出してsupplier表のスキャンを実行し、結果をメモリ上にバッファ
  2. ExecProcNodeを呼び出してcustomer表のスキャンを実行し、結果をメモリ上にバッファ
  3. ExecProcNodeを呼び出してdate1表のスキャンを実行し、結果をメモリ上にバッファ
  4. バッファしたこれらの結果から、GPU-Joinで使用するハッシュ表を共有メモリ上に構築する
  5. GPU-Serviceプロセスにセッションをオープンし、共有メモリ上のハッシュ表をGPUメモリ上にロードする
  6. lineorder表を読み出しつつ、順にGPU-Joinを実行しsupplier、customer、date1の各テーブルと結合処理を行う
  7. GPU-PreAggで集計処理を行い、集計値を結果バッファに書き込む
  8. lineorder表の読み出しが終わると、結果バッファをPostgreSQLバックエンドプロセスに戻す

そこそこの手順を踏んでいる事が分かります。
中でも、(1)と(2)では下位Scanノードの読み出しを行っているのですが、ここではGPU-Scanが使われています。
つまり、supplierやcustomer表をスキャンするために一度GPUへこれらのテーブルのデータをロードしているにも関わらず、一度CPU側へこれを書き戻し、さらにCPUサイクルを消費してこれをバッファにコピーしているわけです。

無駄じゃね?🙄

Star Schema Benchmarkのように、圧倒的に大きな lineorder(この例では876GB)の他は、customer(4.0GB)、supplier(1.3GB)、date1(416kB)程度しかないようなワークロードでは大きな問題になりませんし、そもそもがHash-Joinというアルゴリズムは、サイズが非対称なテーブル同士のJOINを前提としています。少なくともメモリに乗らないとどうしようも無いわけですし。

しかし、バッチ系の処理を中心としてサイズの大きなテーブル同士をJOINするという処理は珍しくないですし、そういった場合に毎度GPUからCPUへとデータを書き戻し、CPU側のメモリを大量に消費するというのは、余分な処理コストと、メモリ消費という2つの点であまりよろしくありません。

Pinned Inner Bufferとは?

PG-StromがSQLワークロードを処理する流れについて簡単に整理してみましょう。
しばしば『GPUってメモリxxGBしか無いのに、xxTBのテーブルを処理できるの?』と聞かれる事があるのですが、その理由はPG-Stromのタスク実行方式にあります。
GPU-JoinやGPU-Scanの場合、読み出すべきテーブル(GPU-Joinの場合はOUTER側のテーブル)を約64MB毎の単位で分割し、この処理単位をGPUに送信してはGPUでこれを評価して、処理結果をCPU側に書き戻します。これを多数のワーカースレッドで並行して行うため、ワーカーAがストレージからの読み出しを行っている間にワーカーBはGPUで処理を実行し、またワーカーCは結果をPostgreSQLバックエンドに書き戻す・・・というように、リソースを極力遊ばせないようにしています。そして、この時に使ったバッファは次の64MBブロックの処理に再利用できるので、テーブルの大きさそれ自体は大きな問題とはなりません。

一方、GPU-PreAggの場合は集計用のバッファをGPU上に置いておき、同様に64MBのブロックからデータを読み出しつつも、この集計バッファを更新した後は処理結果を戻さず、テーブルのスキャンが終わってから一括で処理結果をPostgreSQLバックエンドに戻します。
これは多くの場合、集計クエリが元データを大きく縮減するという特性があるためですが、もちろんカーディナリティの高い重複排除クエリのように読みが外れて事故を起こす事もあり得ます。

さて、ここで使用する結果バッファのデータ形式は、実はGPU-Joinで使用するハッシュ表のデータ形式と同一です。
という事は、GPU-Joinの配下でINNERハッシュ表を作る際に、結果バッファをGPU側に留置しておき(CPU側に戻さず)、それを次のGPU-Joinで使用する事ができれば、余分なデータの移動やCPUによる整形が必要なくなるのではないかと考えました。これがPinned Inner Buffer機構です。

Pinned Inner Bufferの効果を確認するには、もう少し大きなテーブル同士のJOINが必要になるので、ここではTPC-Hのデータセットを利用します。
SF=1000の場合にlineitemは882GB、ordersは205GBですが、ここでは条件句o_orderdate > '1997-05-01'によってInner Bufferのサイズを調整できるようにしています。
TPC-Hのテストクエリそれ自体には、Large Table同士のJOINを行うものはありませんでしたので、ここではデータセットのみ使っています。)

tpch=# set pg_strom.cpu_fallback = off;
SET
tpch=# set pg_strom.pinned_inner_buffer_threshold = '2GB';
SET
tpch=# explain analyze select l_shipmode, o_orderpriority, sum(l_extendedprice)
         from lineitem, orders
        where l_orderkey = o_orderkey
          and o_orderdate > '1997-05-01'
        group by l_shipmode, o_orderpriority;
                                                                            QUERY PLAN

------------------------------------------------------------------------------------------------------------------------------------------------------------------
 Gather  (cost=20248004.34..20248008.55 rows=35 width=59) (actual time=56729.388..56729.515 rows=35 loops=1)
   Workers Planned: 2
   Workers Launched: 2
   ->  Result  (cost=20247004.34..20247005.05 rows=35 width=59) (actual time=56721.837..56721.847 rows=12 loops=3)
         ->  Parallel Custom Scan (GpuPreAgg) on lineitem  (cost=20247004.34..20247004.52 rows=35 width=59) (actual time=56721.835..56721.842 rows=12 loops=3)
               GPU Projection: pgstrom.psum(l_extendedprice), l_shipmode, o_orderpriority
               GPU Join Quals [1]: (l_orderkey = o_orderkey) [plan: 2500102000 -> 480716900, exec: 533500 -> 82944]
               GPU Outer Hash [1]: l_orderkey
               GPU Inner Hash [1]: o_orderkey
               GpuJoin buffer usage: 4096B
               GPU Group Key: l_shipmode, o_orderpriority
               Scan-Engine: GPU-Direct with GPU0; direct=115517748, ntuples=533500
               ->  Parallel Custom Scan (GpuScan) on orders  (cost=100.00..4038700.34 rows=120172200 width=24) (actual time=10926.643..10926.645 rows=-1 loops=3)
                     GPU Projection: o_orderpriority, o_orderkey
                     GPU Pinned Buffer: nitems: 285537432, usage: 23.40GB, total: 27.74GB
                     GPU Scan Quals: (o_orderdate > '1997-05-01'::date) [plan: 1499973000 -> 120172200, exec: 1500000000 -> 285537432]
                     Scan-Engine: GPU-Direct with GPU0; direct=26843246, ntuples=1500000000
 Planning Time: 4.135 ms
 Execution Time: 56730.069 ms
(19 rows)

GpuScan on ordersのパラメータ出力を見てください。
GPU Pinned Buffer: nitems: 285537432, usage: 23.40GB, total: 27.74GBという表示が見えます。これこそがPinned Inner Bufferを使用した事を示す出力で、ここで用いたGPUである NVIDIA A100 (40GB; PCI-E) の3/4近くのサイズを占めるバッファを確保し、これをGPU-Joinで利用したという事になります。

次に、クエリを実行中の I/O の出力の様子を見てみる事にします。I/Oの様子を観察する事で、今、処理がどの辺まで進んでいるのかという事を推測する事ができます。

まず、Pinned Inner Bufferのない従来のGPU-Joinを実行したケースのI/Oの様子です。

面白い事に、クエリの実行開始から25~50秒付近の約25秒、I/Oが停止している時間帯があります。
これがordersテーブルから読み出したデータをInner BufferにCPUで再編するための時間で、2.8億行ともなればそれなりに時間を要している事が分かります。また、前半のordersテーブルの読み出しも10GB/s程度とGpuScanを使っている割には遅く、これはスキャンした結果を受け取るPostgreSQLバックエンドプロセス側でのメモリコピーが間に合っていないと解釈すべきでしょう。

そして、Pinned Inner Bufferを使った場合のI/Oの様子がこちらです。
先ほどのケースにあったような切れ間がなく、概ね20GB/s程度のスループットでLarge Table同士のJOINを処理している事が分かります。

ちなみに、PostgreSQLでのHash-Joinの様子はこのようになっています。
6.0GB/s前後の読み出しが35秒程度持続し、その後は2.8GB/sのスループットでの処理が続いています。
これはおそらく、6.0 * 35 = 210GB なので orders テーブルの読み出しに、そして後半の時間帯はストレージの読み出しよりもむしろ巨大なハッシュ表を使ったHash-Joinの処理がCPUネックとなっていると思われます。

Pinned Inner Buffer の注意点

良い事ばかりに思えるPinned Inner Buffer機能ですが、注意点もあります。

GpuScanの処理結果をGPUメモリに留置するというのが本機能のポイントですが、では、GPUでは処理できなかったデータというのはどうなるのでしょうか?
何を言っているかというと、典型的にはTOAST化された可変長データを参照した場合、などが含まれます。
非常に長い(目安として2kB以上)の可変長データの場合、PostgreSQLはこれを外部の隠しテーブルに格納し、参照があった時にこれを外部テーブルから都度読み出して元の大きさの可変長データを復元します。しかし、GPUに64MB単位のブロックを転送して処理を行っている都合上、そもそもGPUから外部テーブルを参照して処理を行うことなどできません。

そこでPG-StromはCPU-Fallbackという仕組みを用いて、GPUで処理できなかった行はCPUに書き戻して再実行する(= エラーにしない)仕組みがあるのですが、CPUにデータを書き戻すという振る舞いがPinned Inner Bufferと相性が悪く、Pinned Inner Bufferを使用する際にはCPU-Fallbackをoffにセットしなければいけません。
これは副作用のある振る舞いですので、PG-Strom v6.0においてはCPU-Fallbackはデフォルトで有効のままです。したがって、Pinned Inner Bufferを使用するには、明示的にこれらの設定を変更する必要があります。
上記のクエリを実行した際には、以下のような設定を行っています。

tpch=# set pg_strom.cpu_fallback = off;
SET
tpch=# set pg_strom.pinned_inner_buffer_threshold = '2GB';
SET

なお、TOASTメカニズムに関しては、なんと15年前(!)にブログに書いています。
が、実はこの当時と仕組みは変わっていません。凄いですね。
kaigai.hatenablog.com

pg2arrowを並列化する

今回は皆さんが大好きな便利ツール「pg2arrow」のお話です。

PostgreSQLでポータブルな列指向データ形式 Apache Arrow を読み出すには、Arrow_Fdwを利用する事ができます。
PG-StromではGPU-Direct SQLにも対応していますし、列指向データという事もあって、被参照列しかI/Oが発生しない、同じ列のデータが近傍に固まっているという大量データ処理に適した特性を持ってもいます。

また、Apache Arrow形式のファイルを作成するにはPyArrowやPandasなど様々なツールがありますが、我々DB屋としてはPostgreSQLに格納されたトランザクショナルなデータを、分析用にApache Arrow形式として吐き出せるととても嬉しい。そんな時に使えるツールがpg2arrowなのです。

pg2arrowは、PostgreSQLにクエリを投げ、その問合せ結果をApache Arrow形式のファイルとして保存するためのツールで、PG-Stromと同梱して配布されています(ソースコードのかなりの部分を Arrow_Fdw と共用しているためです)。問い合わせに使用するSQLコマンドは、テーブルを単純にダンプするだけではなく、例えばWHERE句で抽出条件を指定したり、複数のテーブルをJOINした結果を含めるといった事も可能です。
GitHubのログを確認したところ、最初のコミットが2019年4月ですので、およそ4年ほど前に設計・開発したツールという事ですね。

一方、pg2arrowには並列動作をサポートしていないという弱点がありました。
そのため、テーブルをダンプするPostgreSQL側のCPUも、データを受け取って書き込むpg2arrow側のCPUも、またその間のネットワークもリソースに余裕があるにも関わらず結構な遊び状態となっており、『例えば1TBくらいのテーブルをArrowに変換してベンチマークを取るぜ!』といった場合でも、夜2:00頃*1にポチっとpg2arrowを走らせて、翌朝結果を確認するという事が常でした。

普通に考えて、並列化すれば大きく高速化するはずです。

データの重複を防ぐために

DBから読み出したデータを別の形式にして保存するだけですので、原理的には、複数のセッションを作成して、それぞれが並列に処理を行えば済む話です。
例えば、あるテーブルをフルダンプしてArrowに変換するとして、各DBクライアントが互いに重なり合わず、しかも漏れの無いような条件で問い合わせを実行すれば良いわけです。典型的には何かのフィールドにハッシュ関数を与え、クライアント数で割った時の剰余がクライアント番号に一致するものだけを取り出せば事足ります。
(これは PostgreSQL のパラレルクエリでも使われているアイデアです)

しかし、考慮すべき点はもう一つあります。
例えば、クライアントAがDBに接続した後、別の誰かが対象テーブルを100行更新し、その後にクライアントB、クライアントCがDBに接続したとします。
この時、更新された100行のうち30行はオレンジ色の、40行は水色の、30行はウグイス色の行だとします。そうすると、結果として生成されたArrowファイルは「部分的に更新された」不整合な状態となってしまいます。

ここでは古典的なテクニックであるスナップショット同期関数を使います。
最初のクライアントAがDBに接続し、トランザクションを開始した後、pg_export_snapshot()を呼び出します。この関数はトランザクションのスナップショットに紐づいたユニークな識別子を返しますが、これを他のセッションでインポートすると、異なるセッションであっても全く同じビューを再現する事ができます。
これはお馴染みpg_dumpの並列ダンプでも用いられており、実際の挙動は以下のようなイメージです。

クライアントA(最初に接続する)

ssbm=# BEGIN READ ONLY;
BEGIN
ssbm=*# SET TRANSACTION ISOLATION LEVEL REPEATABLE READ;
SET
ssbm=*# SELECT pg_catalog.pg_export_snapshot();
 pg_export_snapshot
---------------------
 0000000B-000000B0-1
(1 row)

クライアントB、C、...(ワーカー)

ssbm=# BEGIN READ ONLY;
BEGIN
ssbm=*# SET TRANSACTION ISOLATION LEVEL REPEATABLE READ;
SET
ssbm=*# SET TRANSACTION SNAPSHOT '0000000B-000000B0-1';
SET

pg2arrowを並列に起動してみる。

pg2arrowの並列モードをサポートするために追加されたのは、-n|--num-workers=N_WORKERSオプションと、-k|--parallel-keys=PARALLEL_KEYSオプションの2つです。
これらは互いに排他的で、同じテーブルをスキャンする際に読み出すデータが重ならないよう検索条件を調整するための方法が若干異なってきます。
しかし、内部的な並列処理(ワーカースレッドの挙動)は変わりませんので、対象となるテーブルの設計やデータの特性によって使い分けてください。

$ ./pg2arrow --help
Usage:
  pg2arrow [OPTION] [database] [username]

General options:
  -d, --dbname=DBNAME   Database name to connect to
  -c, --command=COMMAND SQL command to run
  -t, --table=TABLENAME Equivalent to '-c SELECT * FROM TABLENAME'
      (-c and -t are exclusive, either of them must be given)
  -n, --num-workers=N_WORKERS    Enables parallel dump mode.
                        It requires the SQL command contains $(WORKER_ID)
                        and $(N_WORKERS), to be replaced by the numeric
                        worker-id and number of workers.
  -k, --parallel-keys=PARALLEL_KEYS Enables yet another parallel dump.
                        It requires the SQL command contains $(PARALLEL_KEY)
                        to be replaced by the comma separated token in the
                        PARALLEL_KEYS.
      (-n and -k are exclusive, either of them can be give if parallel dump.
       It is user's responsibility to avoid data duplication.)
      --inner-join=SUB_COMMAND
      --outer-join=SUB_COMMAND
  -o, --output=FILENAME result file in Apache Arrow format
      --append=FILENAME result Apache Arrow file to be appended
      (--output and --append are exclusive. If neither of them
       are given, it creates a temporary file.)
  -S, --stat[=COLUMNS] embeds min/max statistics for each record batch
                       COLUMNS is a comma-separated list of the target
                       columns if partially enabled.

Arrow format options:
  -s, --segment-size=SIZE size of record batch for each

Connection options:
  -h, --host=HOSTNAME  database server host
  -p, --port=PORT      database server port
  -u, --user=USERNAME  database user name
  -w, --no-password    never prompt for password
  -W, --password       force password prompt

Other options:
      --dump=FILENAME  dump information of arrow file
      --progress       shows progress of the job
      --set=NAME:VALUE config option to set before SQL execution
      --help           shows this message

Report bugs to <pgstrom@heterodb.com>.
-n|--num-workers=N_WORKERSオプション

-nオプションはシンプルにワーカー数を指定します。
この時、Arrowファイルの元となる問合せSQL文には$(WORKER_ID)$(N_WORKERS)というマクロを埋め込む事ができ、この部分は0から始まるユニークなワーカーIDと、オプションで指定した並列度にそれぞれ置き換えられます。
つまり、-cオプションで指定するSQLをこのように変えれば良いわけです。

オリジナル:

SELECT * FROM lineorder

修正後:

SELECT * FROM lineorder WHERE lo_orderkey % $(N_WORKERS) = $(WORKER_ID)
-k|--parallel-keys=PARALLEL_KEYSオプション

-kオプションはPARALLEL_KEYSに指定したカンマ区切りのトークン毎にワーカースレッドを起動し、そのトークンをそれぞれSQLコマンド中の$(PARALLEL_KEY)に置き換えます。

例を使って説明しましょう。以下のようにパーティション化されたテーブルが存在し、パーティションの子テーブル毎にワーカースレッドを起動して並列にテーブルを読み出したい場合、-kオプションを次のように使います。

ssbm=# \d+
                                                List of relations
 Schema |       Name       |       Type        | Owner  | Persistence | Access method |    Size    | Description
--------+------------------+-------------------+--------+-------------+---------------+------------+-------------
 public | lineorder        | partitioned table | kaigai | permanent   |               | 0 bytes    |
 public | lineorder__p1992 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1993 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1994 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1995 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1996 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1997 | table             | kaigai | permanent   | heap          | 13 GB      |
 public | lineorder__p1998 | table             | kaigai | permanent   | heap          | 7894 MB    |
 public | lineorder__p1999 | table             | kaigai | permanent   | heap          | 8192 bytes |
 public | lineorder_unsort | table             | kaigai | permanent   | heap          | 87 GB      |
(10 rows)

以下の例では、パーティション子テーブルのサフィックスである年号部分を$(PARALLEL_KEY)によって置き換えます。
そうすると、-kオプションで指定したカンマ区切りのキー値ごとにそれぞれワーカーが生成され、結果として、それが各々パーティション子テーブルの条件なしスキャンを実行しています。

$ pg2arrow -d ssbm -c 'SELECT * FROM lineorder__p$(PARALLEL_KEY)' -o /opt/hoge/f_lineorder.arrow -k=1992,1993,1994,1995,1996,1997,1998 --progress
worker:1 SQL=[SELECT * FROM lineorder__p1993]
worker:3 SQL=[SELECT * FROM lineorder__p1995]
worker:2 SQL=[SELECT * FROM lineorder__p1994]
worker:4 SQL=[SELECT * FROM lineorder__p1996]
worker:5 SQL=[SELECT * FROM lineorder__p1997]
worker:6 SQL=[SELECT * FROM lineorder__p1998]
2024-03-31 20:33:58 RecordBatch[0]: offset=1648 length=268436376 (meta=920, body=268435456) nitems=1303083 by worker:0
2024-03-31 20:33:59 RecordBatch[1]: offset=268438024 length=268436376 (meta=920, body=268435456) nitems=1303083 by worker:3
          :
          :
worker:0 merged pending results by worker:4
2024-03-31 20:38:15 RecordBatch[460]: offset=123480734608 length=127664216 (meta=920, body=127663296) nitems=619722 by worker:0
Total elapsed time: 00:04:26

ワーカーの動作について

ここで、大雑把な処理の流れにも触れておく事にします。
Pg2Arrowが並列モードで起動すると、mainスレッドであるworker-0がスキーマ定義を作成するなどの初期設定を行い、その後、他のワーカーを順次起動していきます。
各ワーカーはそれぞれPostgreSQLに接続し、それぞれ約256MBのバッファ((-sオプションで変更可))が埋まるたびに出力先のArrowファイルへと書き込みを行います。
Arrowファイルは内部がRecord Batchと呼ばれるブロックに分割されており、ファイルポインタを進める部分さえアトミックに行ってしまえば、以降の書込み処理はマルチスレッドが複数の別個の領域に独立して書き込む事ができます。そのため、シーケンシャルに実行しなければならないクリティカルセクションは最小限に抑えられています。

また、クエリの最後まで読み出したにも関わらず256MBのバッファを埋められなかった場合は、隣接スレッドのバッファにマージされ、最終的には worker-0 のバッファにマージされますので、並列ワーカーの数を増やしたとしても Record Batch の大きさが極端に小さい(= PG-Stromでの実行効率が低下する)データが作られるわけではありません。


生成した Arrow ファイルをFDW経由で参照する

それでは、生成した Arrow ファイルを参照してみる事にします。
個々のテーブル定義(カラム名、データ型)を指定するのは面倒ですので、IMPORT FOREIGN SCHEMA文を使用するのがお勧めです。

ssbm=# import foreign schema f_lineorder from server arrow_fdw into public options (file '/opt/hoge/f_lineorder.arrow');
IMPORT FOREIGN SCHEMA
ssbm=# \d f_lineorder
                        Foreign table "public.f_lineorder"
       Column       |     Type      | Collation | Nullable | Default | FDW options
--------------------+---------------+-----------+----------+---------+-------------
 lo_orderkey        | numeric       |           |          |         |
 lo_linenumber      | integer       |           |          |         |
 lo_custkey         | numeric       |           |          |         |
 lo_partkey         | integer       |           |          |         |
 lo_suppkey         | numeric       |           |          |         |
 lo_orderdate       | integer       |           |          |         |
 lo_orderpriority   | character(15) |           |          |         |
 lo_shippriority    | character(1)  |           |          |         |
 lo_quantity        | numeric       |           |          |         |
 lo_extendedprice   | numeric       |           |          |         |
 lo_ordertotalprice | numeric       |           |          |         |
 lo_discount        | numeric       |           |          |         |
 lo_revenue         | numeric       |           |          |         |
 lo_supplycost      | numeric       |           |          |         |
 lo_tax             | numeric       |           |          |         |
 lo_commit_date     | character(8)  |           |          |         |
 lo_shipmode        | character(10) |           |          |         |
Server: arrow_fdw
FDW options: (file '/opt/hoge/f_lineorder.arrow')

試しに、SSBMのQ1_2を走らせてみます。

ssbm=# select sum(lo_extendedprice*lo_discount) as revenue
from f_lineorder, date1
where lo_orderdate = d_datekey
  and d_yearmonthnum = 199401
  and lo_discount between 4 and 6
  and lo_quantity between 26 and 35;
    revenue
---------------
 9624332170119
(1 row)

もちろん、Arrowファイルの元となったlineorderテーブルを使っても同じ値が返ってきます。

ssbm=# select sum(lo_extendedprice*lo_discount) as revenue
from lineorder, date1
where lo_orderdate = d_datekey
  and d_yearmonthnum = 199401
  and lo_discount between 4 and 6
  and lo_quantity between 26 and 35;
    revenue
---------------
 9624332170119
(1 row)

並列Pg2Arrowのパフォーマンス

最後に、もちろん気になる並列Pg2Arrowのパフォーマンスを見てみる事にします。

測定環境のスペックは以下の通り。

以下のオプションで pg2arrow を起動し、87GBのlineorderテーブルを全てArrow形式に変換するまでの時間を計測しました。

$ pg2arrow -d ssbm -c 'SELECT * FROM lineorder_unsort WHERE lo_orderkey % $(N_WORKERS) = $(WORKER_ID)' -o /opt/hoge/f_lineorder.arrow -n N_WORKERS --progress

まず緑の縦棒が並列実行なしのpg2arrowで87GBのlineorderテーブルをArrowに変換した時のもので、1,520秒、おおむね25分ちょい要しています。
これを-nオプションで並列数を増やした場合、およそ400秒(6分40秒)を少し下回った辺りで頭打ちとなっているようです。

これは検索条件(lo_orderkey % $(N_WORKERS) = $(WORKER_ID))で重複行を排除しているため、クライアント数が増加するにしたがって、同時に lineorder テーブルを重複してスキャンしなければならなくなり、トータルのバッファからの読出し負荷が増えてしまったためと言えるでしょう。

一方、パーティション子テーブルの名前を一部を-kオプションで置き換え、それぞれのクライアントが完全に独立した領域をスキャンする事になったパターンでは、並列度が7である(しかもlineorder__p1998の容量は他の子テーブルの半分なので、実質6.5並列)にも関わらず、266秒と並列実行なしのパターンに比べて5.7倍の実行時間を記録しています。

この事から、並列Pg2Arrowのパフォーマンス向上を享受するには、以下の点に注意を払う必要があるでしょう。

  • 単純に並列度を増やすだけでも効果はあるが、同じ領域を重複してスキャンする処理が増えると、効果は限定的になりがち。
  • PostgreSQLがスキャンする量を減らすような検索条件の与え方が望ましい。パーティションなどで分割されていたら理想的。

*1:もっと早く寝ろw

PG-Strom v5.0

ずいぶんご無沙汰のブログ記事となりました。

今回は、設計を一新して速く、頑強になった PG-Strom v5.0 をご紹介します。

なぜ再設計が必要だったのか?

前バージョンの PG-Strom v3.x シリーズの基本的な設計は、2018年のPG-Strom v2.0の頃から大きく変わっていません。
当時の最新GPUモデルは Volta 世代(TESLA V100)で、CUDAのバージョンは9.2ですから、かなりの大昔という事はお分かり頂けると思います。

この頃、PG-Stromの開発において最優先すべき課題は、先ず実用となるバージョンをリリースする事でした。(※ HeteroDB社の創業は2017年7月です)
クエリの処理速度を高速化する事は当然なのですが、それ以上に、まだPG-Stromの内部インフラも十分に枯れていない中で、クラッシュせずに走り切る事や、バグがあったとしても容易に原因箇所を特定できる事が優先であったのです。また、GPU側でSQLを実行するデバイスコードにしても、様々な実装方式にトライしてその中で最良を選択するというよりも、先ずは出たとこ勝負で『動くモノ』を優先するという状況でした。

ただその後、数年を経て、明らかになってきた問題が複数あり、これらはどこかのタイミングで大規模なリファクタリングを行わざるを得ないと考えていました。例えば以下のような問題点です。

問題①:CUDA Contextが消費するリソース

ご存知のようにPostgreSQLはマルチプロセスで動作します。クライアントからの接続が発生するたびにバックエンドプロセスをfork(2)し、そのプロセスがSQL処理の大半を担います。また、サイズの大きなテーブルをスキャンする時など、一時的なワーカープロセスを起動してSQL処理を並列に実行する事もあります。
PG-Stromの追加した実行計画(GpuScan、GpuJoin、GpuPreAgg)が採用された場合、これらはGPUを使用するためにCUDA Contextというものを作成します。しかしCUDA Contextを作成すると、それ自体がGPUのリソース(デバイスメモリ数百MB~)を消費してしまうため、PostgreSQLへの同時接続数が増加すると、ワーキングに利用できるメモリがほとんど残らない事になります。

問題②:複雑すぎるGPUコード自動生成ロジック

2012年にPG-Stromの最初のプロトタイプを作成した時から、PG-StromはSQLとして与えられたScan条件式(WHERE句)やJoin結合条件(ON~句)から自動的にCUDA C++ソースコードを生成し、それを実行時コンパイルしてGPU用のネイティブバイナリを生成していました。
しかし、様々な状況に対応してSQLからCUDA C++用のソースコードを生成するロジックは非常に複雑で、例えば、GpuJoin用のソースコード生成を含むsrc/gpujoin.cは8500行近くの規模があり、ソースコードを保守する上でかなり悩ましい問題を抱えていました。(要はスパゲッティ)

問題③:最低限必要な300ms

PG-Strom v3.x以前は各バックエンドプロセスがCUDA Contextを作成し、GPUのメモリ割り当てやタスク(GPU Kernel)の投入を行っていました。
このCUDA Contextの作成には実は少し時間がかかり、およそ100~150ms程度の遅延が不可避でした。また、CUDA C++ソースコードを作成し、これをGPU向けのバイナリにコンパイルする際にも、最低で200ms程度の時間がかかっていました(もちろん処理の複雑さによります)。
何十秒もかかる処理ならともかく、数十ms程度の応答速度を要求されるクエリでこのペナルティは割と厳しいものがあります。

問題④:NVIDIA GPU以外への拡張性

これは現時点では可能性の話にすぎませんが、例えば、PG-Stromの仕組みをComputational Storage Drive (CSD)に実装する事ができれば、ストレージ側のプロセッサでScan、Join、GroupByといったSQL処理の一部を実行したり、Projection処理を行う事で被参照列のみをホストに返すという列指向データ構造に近い事ができるはずです。しかし、PG-StromがCUDA C++を前提としたソースコード生成に注力していた場合、CSDで実行可能な処理の自動生成部分を二重に持つ事となり、開発効率以上にソフトウェア品質的に悩ましい問題を抱える事となりそうです。

PG-Strom v5.0のアーキテクチャ

これらの問題を解決するため、PG-Strom v5.0ではまずPostgreSQLの各バックエンドプロセスがCUDA Contextを持つ構造を廃止しました。
代わりに、常駐プロセスであるPG-Strom GPU-Serviceだけが唯一GPUと相対してリソースの管理やタスクの投入を行います。
PostgreSQLのバックエンドプロセス(で動作するPG-StromのCustomScanハンドラ)は、プロセス間通信を通じてGPU-Serviceにリクエストを送出し、その処理結果を待つだけです。GPU-Serviceはマルチスレッドで動作し、pg_strom.max_async_tasksを上限として並列にタスクを処理する事ができます。

以下の模式図をご覧ください。
PostgreSQLバックエンドプロセスが個々にGPUを管理する場合と、PG-Strom GPU-ServiceだけがGPUを管理し、他のプロセスはGPU-Serviceにリクエストを送出するモデルでは、CUDAによって消費されるGPUバイスメモリの量が段違い(特にクライアント数が多い場合)である事がよく分かると思います。

さらにこの構造は、CUDA Contextを初期化するとき(cuCtxCreate())の100ms~150ms遅延の問題も解決します。なぜなら、CUDA ContextはGPU-Serviceの起動時に既に作成済みで、それに比べればUNIXドメインソケットを通じてGPU-Serviceとの間にコネクションを確立する処理時間など微々たるものにしかすぎないからです。

CUDA C++ネイティブコードから、疑似命令コードへ

もう一つ。これまでSQLワークロードをGPUで実行するためにCUDA C++ソースコードを生成し、これを実行時コンパイラ(NVRTC:NVIDIA Run Time Compiler)が最適化、実行時バイナリの生成というステップを踏んでいましたが、GpuJoinなどでCUDA C++ソースコードを生成するためのロジックが複雑になりすぎた事から、PG-Strom v5.0では疑似命令コードを生成するようになりました。

以下の実行計画をご覧ください。
これは様々な条件で絞り込みを行ったdate1テーブルとlineorderテーブルをJOINし、lo_extendedprice*lo_discountの結果を集計するクエリの実行計画です。
VERBOSEオプションを付加すると、GpuPreAggプランの下の方に『xxx OpCode』というものが出力されています。(※ VERBOSEオプション抜きだとここまで賑やかなEXPLAIN出力にはなりません)
このOpCodeというのは、GPU上で実行する演算子や列参照の手順をバイナリ形式でパッキングしたもので、EXPLAINの出力では可読な形式に直したものを出力しています。

従来は、ここに生成したCUDA C++ソースコードのファイル名が出力されていました。これをNVRTCに渡して実行時コンパイルし、GPU用のバイナリを生成するわけです。

=# explain verbose
select sum(lo_extendedprice*lo_discount) as revenue
from lineorder,date1
where lo_orderdate = d_datekey
and d_year = 1993
and lo_discount between 1 and 3
and lo_quantity < 25;
                                                                    QUERY PLAN
-----------------------------------------------------------------------------------------------------------------
 Aggregate  (cost=12713126.00..12713126.01 rows=1 width=32)
   Output: pgstrom.sum_fp_num((pgstrom.psum(((lineorder.lo_extendedprice * lineorder.lo_discount))::double precision)))
   ->  Custom Scan (GpuPreAgg) on public.lineorder  (cost=12713125.99..12713126.00 rows=1 width=32)
         Output: (pgstrom.psum(((lineorder.lo_extendedprice * lineorder.lo_discount))::double precision))
         GPU Projection: pgstrom.psum(((lineorder.lo_extendedprice * lineorder.lo_discount))::double precision)
         GPU Scan Quals: ((lineorder.lo_discount >= '1'::numeric) AND (lineorder.lo_discount <= '3'::numeric) AND (lineorder.lo_quantity < '25'::numeric)) [rows: 2400065000 -> 315657500]
         GPU Join Quals [1]: (date1.d_datekey = lineorder.lo_orderdate) ... [nrows: 315657500 -> 45076290]
         GPU Outer Hash [1]: lineorder.lo_orderdate
         GPU Inner Hash [1]: date1.d_datekey
         GPU-Direct SQL: enabled (GPU-0)
         KVars-Slot: <slot=0, type='numeric', expr='lineorder.lo_discount'>, <slot=1, type='numeric', expr='lineorder.lo_quantity'>, <slot=2, type='float8', expr='(lineorder.lo_extendedprice * lineorder.lo_discount)'>, <slot=3, type='numeric', expr='lineorder.lo_extendedprice'>, <slot=4, type='int4', expr='date1.d_datekey'>, <slot=5, type='int4', expr='lineorder.lo_orderdate'>
         KVecs-Buffer: nbytes: 192512, ndims: 3, items=[kvec0=<0x0000-dfff, type='numeric', expr='lo_discount'>, kvec1=<0xe000-1bfff, type='numeric', expr='lo_quantity'>, kvec2=<0x1c000-29fff, type='numeric', expr='lo_extendedprice'>, kvec3=<0x2a000-2c7ff, type='int4', expr='d_datekey'>, kvec4=<0x2c800-2efff, type='int4', expr='lo_orderdate'>]
         LoadVars OpCode: {Packed items[0]={LoadVars(depth=0): kvars=[<slot=5, type='int4' resno=6(lo_orderdate)>, <slot=1, type='numeric' resno=9(lo_quantity)>, <slot=3, type='numeric' resno=10(lo_extendedprice)>, <slot=0, type='numeric' resno=12(lo_discount)>]}, items[1]={LoadVars(depth=1): kvars=[<slot=4, type='int4' resno=1(d_datekey)>]}}
         MoveVars OpCode: {Packed items[0]={MoveVars(depth=0): items=[<slot=0, offset=0x0000-dfff, type='numeric', expr='lo_discount'>, <slot=3, offset=0x1c000-29fff, type='numeric', expr='lo_extendedprice'>, <slot=5, offset=0x2c800-2efff, type='int4', expr='lo_orderdate'>]}}, items[1]={MoveVars(depth=1): items=[<offset=0x0000-dfff, type='numeric', expr='lo_discount'>, <offset=0x1c000-29fff, type='numeric', expr='lo_extendedprice'>]}}}
         Scan Quals OpCode: {Bool::AND args=[{Func(bool)::numeric_ge args=[{Var(numeric): slot=0, expr='lo_discount'}, {Const(numeric): value='1'}]}, {Func(bool)::numeric_le args=[{Var(numeric): slot=0, expr='lo_discount'}, {Const(numeric): value='3'}]}, {Func(bool)::numeric_lt args=[{Var(numeric): slot=1, expr='lo_quantity'}, {Const(numeric): value='25'}]}]}
         Join Quals OpCode: {Packed items[1]={JoinQuals:  {Func(bool)::int4eq args=[{Var(int4): slot=4, expr='d_datekey'}, {Var(int4): kvec=0x2c800-2f000, expr='lo_orderdate'}]}}}
         Join HashValue OpCode: {Packed items[1]={HashValue arg={Var(int4): kvec=0x2c800-2f000, expr='lo_orderdate'}}}
         Partial Aggregation OpCode: {AggFuncs <psum::fp[slot=2, expr='(lo_extendedprice * lo_discount)']> arg={SaveExpr: <slot=2, type='float8'> arg={Func(float8)::float8 arg={Func(numeric)::numeric_mul args=[{Var(numeric): kvec=0x1c000-
2a000, expr='lo_extendedprice'}, {Var(numeric): kvec=0x0000-e000, expr='lo_discount'}]}}}}
         Partial Function BufSz: 16
         ->  Seq Scan on public.date1  (cost=0.00..78.95 rows=365 width=4)
               Output: date1.d_datekey
               Filter: (date1.d_year = 1993)
(22 rows)

しかし、PG-Strom v3.xが自動生成するCUDA C++コードは、ライブラリ部分を予めビルドしておく方式に切り替えた事から、実質的にはSQL文が与えられるたびに変化する制御構造をソースコード自動生成という形で吸収していたとも言えます。
そこで、この制御構造自体をGPU側へ持ち込めば(+条件分岐に相当する部分は予め関数ポインタをセットするなどして実行コストを抑える)、実行時コンパイルの手間を省けると考えたわけです。

kaigai.hatenablog.com

実際に、比較的規模の小さなテーブル(800万件)の集計処理を PG-Strom v3.5 と PG-Strom v5.0 で比較してみます。

  • PG-Strom v3.5
ssbm=# select count(*) from lineorder_8m where lo_orderpriority = '2-HIGH';
  count
---------
 1604233
(1 row)

Time: 1132.471 ms (00:01.132)
  • PG-Strom v5.0
=# select count(*) from lineorder_8m where lo_orderpriority = '2-HIGH';
  count
---------
 1604233
(1 row)

Time: 114.781 ms

全体で数十秒~を要するクエリであれば初期セットアップの時間差は大きく影響しませんが、比較的小さなデータセットであれば、クエリの実行時間を頑張って速くしてもある一定の限度以上には高速化できないという問題がありました。しかし、v5.0ではGPU-Serviceが既にCUDA Contextを初期化している上、コードのコンパイル&最適化も不要であるため、比較的小さなテーブルであってもGPU処理の恩恵を得やすいというメリットがあります。

GPU-Serviceのマルチスレッド化とパフォーマンス改善

v5.0での大規模なリファクタリングによって、GPUを管理するのはGPU-Serviceプロセス一個だけに絞られ、さらにGPU-ServiceはマルチスレッドによりPostgreSQLバックエンドプロセスからのリクエストを次々と捌いていきます。これはGPUやCUDAのレイヤから見ると大きな変化で、並列に動作するPostgreSQLバックエンドからの要求を処理するたびにCUDA Contextを切り替える必要がなくなり、GPU Kernelを起動する際のスループットが向上します。

これは処理速度がシビアな状況で引っかかってくる事があり、例えばこれは同じ構成のサーバ*1上でStar Schema Benchmark (SSBM)を実行した場合、このSSDのSeqRead速度は6500MB/sですので、理論上、4本束ねた場合は26,000MB/sまでの読出しスループットを発揮できるはずです。
しかし、v3.5の結果を見ると20GB/s程度で性能値が頭打ちになっている一方、v5.0では24GB/s程度まで処理性能が伸びている事が分かります。かなり多くの部分で修正が加えられているため、これだけが高速化の要因というわけではないでしょうが、v5.0になりGPUへタスクを放り込むスケジューリングがより洗練されるようになってきたという事が分かります。


まとめ

PG-Stromにおけるこれら内部アーキテクチャの一新は、安定性・保守性を大きく高めると共に、GPU-Direct SQLでよりハードウェア理論速度に近いパフォーマンスを発揮し、さらにCUDA Contextの生成やCUDA C++コードのコンパイルに要する時間の削減効果で、とりわけ比較的小さなデータセット(~20GB程度)であってもGPU利用の効果が実感できるようになりました。

これらPG-Strom v5.0の特徴、修正点については、明日(3/15)のセミナーでお話しさせていただきますので、ぜひこちらも併せてご参加いただければと思います。

bakusokudb.connpass.com

*1:CPU: AMD EPYC 7402P (24C, 2.8GHz)、RAM: 128GB、GPU: NVIDIA A100 [40GB, PCI-E]、SSD: Intel SSD D7-P5510 [U.2, 3.84TB]