シャーディングされたPostgresクエリのライフサイクル
Postgresが千台のサーバーにシャーディングされた後、1つのSELECTが認証、wire protocol、分散クエリプランナーをどう通過するか、NekiルーターはGoでSCRAM-SHA-256を再実装しSimpleとExtendedの両プロトコルをサポート。
日本語
コピー

Postgres を使ったことがあり、「なんてシンプルなソフトウェアなんだ、どの部分も完全に理解できる」と思ったなら、それはクエリプランナーを見ていないだけだ。ましてや、そのデータベースを千台のサーバーにシャーディングすることを考えてみてほしい。クエリはどう処理する?Postgres の認証機構、wire protocol、パーサーを再現し、シャードを意識した分散クエリプランナーを持ち、あらゆるサーバー障害シナリオを優雅に処理し、Postgres のコネクションごとにプロセスを立てるアーキテクチャを回避するコネクションプールも必要だ。簡単だろう?あるひとつの Postgres クエリを追いながら、この優雅で精巧かつ複雑な分散システムの各レイヤーを順に見ていこう。そうすれば、数千台のサーバーにまたがる大規模な Postgres シャーディング構成を、あたかも一台の Postgres サーバーのように見せるために何が必要なのかが理解できるはずだ。実際のデータベースには何百ものテーブルがあるだろうが、この例ではスキーマをシンプルに保つ。テーブルは2つ、4つのシャードに分散している。
CREATE TABLE customers (
id BIGINT PRIMARY KEY,
name TEXT NOT NULL,
email TEXT NOT NULL,
country TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE orders (
id BIGINT PRIMARY KEY,
customer_id BIGINT NOT NULL,
total NUMERIC(12,2) NOT NULL,
status TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
追いかけるクエリもシンプルだ。SELECT 文がひとつ、最近の注文すべてについて顧客名、注文金額、注文日を取得する。
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
このクエリを実行する前に、シャーディングについて理解しておくべき重要なことがある。シャーディングの意義は、データベースを単一サーバーの限界から解き放つことだ。つまり、ここから先は分散システムである。行データは異なるサーバーにまたがりうるので、データをどこに置き、どうやって見つけるかのルールが要る。この例では、customers.id を customers のシャードキーに、orders.id を orders のシャードキーに選ぶ。これらの主キーは自然な出発点だが、少し素朴でもあり、それはすぐに分かる。customers や orders の行を保存するたびに、システムは id を取り出し、ハッシュを計算し、そのハッシュで保存先を決める。顧客と注文が多くのサーバーに散らばっていることに注目してほしい。これがシャーディングされたデータベースをほぼ無限にスケールさせる一方で、興味深いエンジニアリング上の難題をいくつも持ち込む。Follow Maja ボタンで、ある顧客とその注文がどうシャーディングされるか見てみよう。目ざとい人なら、両者が同じシャードにないことに気づくだろう。これは覚えておいてほしい。すぐに効いてくる。それでは Neki のシャーディングデータベースを巡る旅を始めよう。
認証
アプリケーションがクエリを送るには、まず認証を通過しなければならない。シャーディングデータベースの中には、ルーティングをアプリケーション層に任せるものもある。アプリが適切な Postgres インスタンスを自分で選び、接続を張る。我々のルーターは、その複雑さを隠す。あなたのアプリケーションにとって、ルーターは_まさに_ Postgres であり、認証、接続、そしてクエリの各サーバーへの振り分けを担う。そのために、SCRAM-SHA-256 を含む Postgres の認証フローを Go で完全に実装した。TCP の確立、SSL ネゴシエーション、TLS ハンドシェイクの後、クライアントは startup メッセージを送る。そこには Postgres のロール、データベース名、要求されたセッション設定が含まれる。ルーターは自身の認証ルールを確認し、SCRAM の交換が始まる。
ルーターが challenge を送る。ドライバはパスワードと challenge から proof を計算し、ルーターはそれを保存済みの verifier と照合する。続いてルーターはドライバが検証するための signature を返す。こうしてルーターは、クライアントが正しいパスワードを知っていることを確認できる。パスワードそのものをネットワーク上に流すことなく。
最後にルーターはデータベースへのアクセス権限とアプリケーションのセッション設定を確認し、startup レスポンスを送る。ReadyForQuery がクライアントにクエリを送ってよいと伝える。
プロトコル
安全な接続が確立すると、アプリケーションはクエリを Postgres ドライバに渡せる。あまり知られていないが、Postgres の wire protocol にはクエリを送る方法が2つある。Simple と Extended だ。さらに、Postgres ドライバによって、あるいは同じドライバでも設定によって、同じクエリの送り方が変わりうることも知られていないかもしれない。ルーターの目標はあらゆる Postgres クライアントと互換であることなので、両方をサポートする。
Postgres Simple プロトコルと Postgres Extended プロトコル
simple プロトコルでは、クライアントは SQL 文全体をタイムスタンプとともに1つの Query メッセージで送る。
Query
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= '2026-09-09 00:00:00+00';
extended プロトコルでは、クエリは5つの部分に分けて渡される。
Parse
statement: ""
query: SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
Bind
statement: ""
portal: ""
parameters: ["2026-09-09 00:00:00+00"]
Describe
portal: ""
Execute
portal: ""
Sync
- Parse はルーターに SQL の準備を要求する。この時点で
$1はまだそのままの位置にある。これにより無名の prepared statement が作られる。アプリケーションが明示的にPREPAREを実行していなくてもだ。 - Bind はタイムスタンプを渡し、この実行のための portal を作る。portal はクエリの実行状態を追跡し、prepared statement をそのパラメータ値と要求された結果フォーマットに結びつける。
- Describe は、この portal がどの列をどのフォーマットで返すかを尋ねる。
- Execute はルーターにその portal を実行させる。
- Sync はこのバッチの終わりを示す。
psql のようなコマンドラインから来るアドホックなクエリにとって、simple プロトコルは シンプル なやり取りを提供する。完全な SQL を送り、結果を受け取る。1回のリクエストで複数の SQL 文を受け付けることもできる。一方 extended プロトコルは、Postgres ドライバが準備、パラメータのバインド、実行を別々に制御できるようにする。アプリケーションが同じ文を異なる値で繰り返し実行するとき、この分離が役に立つ。準備済みの文を再利用すれば、SQL を毎回パースし、その中の名前や型を毎回解決せずに済む。我々の router は本物の Postgres プロセスではないので、このプロトコルの振る舞いを実装するにあたって独自の最適化もいくつか加えている。ただそれは別の記事の話題だ。アプリケーションから見れば、我々の router は Postgres そのものだ。
解析
単純なプロトコルであれ拡張プロトコルであれ、最終的に手に入るのは処理必須の SQL 文字列だ。Postgres と直接やり取りする場合、Postgres のパーサーがまずクエリを字句解析して抽象構文木(AST)を構築し、それを Postgres のクエリプランナーに渡す。AST はクエリの構成要素を表す。どのテーブルを参照し、どう結合し、どんなフィルタ条件を課し、どのカラムを返すのか。router 側でも同じ解析と木の構築を行わなければならず、しかも Postgres の文法と正確に一致させる必要がある。受け入れる構文も、拒否する構文も同じでなければならない。これを約 18,000 行の Go で実装し、厳格な要件とテストを課すことで、Postgres 自身の C 実装に劣らず、場合によっては上回る性能を確保した。ちょうどコンパイラのパーサーのように、ソーステキストを抽象構文木へと変え、router がクエリの実行方法を検討する際に検証・書き換えできる構造を与える。router が SQL を解析する必要すらないケースもある。性能のために、router は無駄な処理を可能な限り避けることを目指す。まず再利用可能なキャッシュ済みプランがあるかを確認し、キャッシュにヒットすれば解析と計画をスキップできる。だが今回はヒットしなかった。router は解析せざるを得ず、その AST がクエリプランナーに渡される。
計画
ここまでで多大な作業をこなしてきたが、まだ SQL クエリを抽象構文木に変えただけだ。シャーディングはまだ考慮されておらず、Postgres にも触れておらず、「どうすればこのクエリを高速に実行できるか」という問いにも答えていない。
シャーディングされたデータベースシステムの良し悪しは、クエリプランナーがどれだけ堅牢かで決まる。腰を据えてほしい。Neki にできるだけ働かせないために我々がどれだけ手を尽くしたか、すぐに分かるだろう。とことん怠け者なのだ。
ここまで追ってきた元のクエリをもう一度見てみよう。
SELECT customers.name, orders.total, orders.created_at
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1;
Postgres ノードが 1 つだけだったら、これを実行するには何が必要か? シャーディングもプロキシもクエリルーターもない(今のところは)。Postgres には独自のクエリプランナーがあり、パーサーから AST を受け取ると、次のことを決めなければならない。
- 内部統計を使って、
created_at >= $1に一致する注文がいくつあるか見積もる。 - 各テーブルの読み取り方法を決める。シーケンシャルスキャン、インデックススキャン、ビットマップスキャン、インデックスオンリースキャン。
- 結合の方向を選び、ネステッドループ、ハッシュ結合、マージ結合のどれを使うか決める。
orders.created_atのインデックスが注文のフィルタリングに役立つか確認する。
これでもこの単純なクエリに必要な判断の一部にすぎない。クエリが複雑になれば、プランナーが考慮すべきことは増える。結合、集約、サブクエリ、ソートが増えていく。
忘れないでほしい。これはシャーディングされたデータベースで、各テーブルの行は各シャードに散らばっている。ここでは 4 シャードだが、この節で述べる原理は、4 シャードに 1 TB でも、400 シャードに 1 PB でも成り立つ。
各行がどのシャードに落ちるかを決めるために、あるカラムをシャードキーとして指定する。以前、customers には customers.id を、orders には orders.id を選んだ。前の図をもう一度見てほしい。この決定はすぐに厄介な問題になる。
この配置では、ある顧客の行は 1 つのシャードに落ちるが、その顧客の注文は多くのシャードに散らばる。さて、どう結合すればいいのか?
シャーディングされたクエリプランナーは、任意の SQL クエリを受け取り、複数の Postgres ノードにまたがる実行のために最適なプランを構築できなければならない。
プランの構築
クエリプランナーの成果物もまた木だ。プランナーの役割は、クエリ要求を記述した AST をクエリプランへ変換することであり、クエリプランは結果を算出するためにクラスタが完了しなければならない操作列のすべてを記述する。そのために、router 内のプランナーは 2 つのものを同時に把握していなければならない。(a) 完全な Postgres schema と、(b) データが各シャード間でどう分布しているかの規則だ。(a) はデータベースの_権威シャード_(authoritative shard)から取得する。このシャードは、このデータベース schema メタデータの供給元として指定されている。router はここからテーブル名、カラム、型を取得する。クエリのたびに権威シャードへ問い合わせるのではなく、キャッシュを保持し、定期的に権威シャードと通信して更新を受け取る。(b) はデータトポロジーから取得する。このトポロジーは etcd に集中管理され、各 router ノードにキャッシュされる。これらの情報を頼りに、router はクエリ内のカラムとその型を解析する。この手順を意味解析と呼ぶ。分散プランニングを回避できる近道がないかも探るが、あいにく我々のクエリにはその運がない。
どう join するか?
次は AST をクエリプラン木に変える。最初の一歩は粗いプラン木に変換し、その後いくつかの段階で詳細化する。元のクエリから生成される基本プランはおおよそこうなる。これは単純だ。JoinCluster が 2 つのテーブルとその inner join 条件をまとめ上げる。Inner join は並べ替えが可能で、これが同じクエリに対してプランナーに異なる対処法を与える。OutputClauses は行を取得した後に何をするかを記述する。ここでは 3 カラムを返すだけでよい。他のクエリではソート、グループ化、集約、LIMIT、さらには DISTINCT が必要かもしれない。プランナーはこれらの行をどう join するかも決めなければならず、ここがより難しい部分だ。
シャードをまたぐ join
プランナーはデータトポロジーを使って、クエリのどの部分をシャード上で実行できるか判断する。先ほどは示さなかったが、このルーティングマッピングはデータトポロジーの JSON ファイルで指定している。2 つのテーブルについて、それぞれ自身の id のハッシュでシャーディングすると、マッピングはおおよそ次のようになる(データベース名は postgres)。
{
"authoritative_shard_group": "metadata",
"shard_indexes": {
"xxhash_id": {
"type": "xxhash",
"columns": ["id"]
}
},
"shard_groups": [
{
"uid": "metadata",
"key_ranges": [{ "shard_uid": "shard1" }]
},
{
"uid": "main_shards",
"default_shard_index": "xxhash_id",
"key_ranges": [
{ "shard_uid": "shard1", "end": "40" },
{ "shard_uid": "shard2", "start": "40", "end": "80" },
{ "shard_uid": "shard3", "start": "80", "end": "c0" },
{ "shard_uid": "shard4", "start": "c0" }
]
}
],
"databases": {
"postgres": {
"schemas": {
"public": {
"tables": {
"customers": { "shard_group": "main_shards" },
"orders": { "shard_group": "main_shards" }
}
}
}
}
}
}
この設定は router にこう伝える。customers と orders の両テーブルを main_shards グループの 4 つのシャード(shard)に分散させ、それぞれ id カラムのハッシュ値でシャーディングする。新しい行を格納するとき、router は次のようにする。
- その行の
idを取り出す。 - それを決定的な
xxhash関数に入力する。 - 得られたハッシュ値をもとに、あらかじめ決めた範囲からシャードを選ぶ。
4 つのシャードはそれぞれ異なるハッシュ値の区間を担当している。前に触れたとおり、シャードキーをこう選ぶと話がややこしくなる。最初のしわ寄せがここで出る。ある顧客の注文は 4 つのシャードすべてに散らばりうるのに、顧客の行そのものはどれか 1 つのシャードにしかない。両者を結合するには、ルーター側で一致する行を突き合わせるしかない。とはいえ我々の planner は賢い。JoinCluster ノードを更新し、行を取得して結合する計画へと作り変えていく。このクエリでは、customers と orders の処理を切り分け、各シャードに適用できるフィルタ条件を洗い出し、ルーターがどう結果を突き合わせるかを決める必要がある。まずテーブルごとに Route の計画ノードを作る。各 route の内部処理は、Postgres に渡す 1 本の SQL クエリになる。route は対象シャードの選び方も指定するので、1 つのシャードに送ることも複数に送ることもできる。customers と orders の両テーブルから必要な行を取り出すには scatter-gather クエリを送ることになる、とここまでで分かった。だが、この結合にはいくつもやり方がある。顧客を先に取るか、注文を先に取るか。注文の絞り込みはシャード側でやるのか、ルーター側でやるのか。結合はネストループでやるのか、ハッシュテーブルでやるのか。これが planning の次の段階だ。
複数の Postgres インスタンスをまたぐ結合
JOIN orders ON orders.customer_id = customers.id を実現するには、ルーターが条件に合う注文と対応する顧客を突き合わせなければならない。方法はいくつかある。ネストループ結合(nested-loop join)は、データベースの内部実装に馴染みのある人なら聞き覚えがあるだろう。そうでない人向けに言えば、発想は単純だ。片側の入力から行を 1 つ取り、その行ごとにもう片側の入力から一致するものを探す。ルーターはまず各シャードにクエリを送り、customers の行をすべて集める。行が返ってくるたびに、ルーターは顧客を 1 人ずつ見ていく。顧客 1 人につき、全シャードへ scatter-gather クエリを 1 回投げ、その顧客の注文をすべて集める。条件に合う注文が 1 件もない顧客であっても、それを確認するだけで 4 回のシャードクエリが飛ぶ。router
| id | name |
|---|---|
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
SELECT id, nameFROM customers;
| name | total | created_at |
|---|---|---|
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
分片 1customers
| 2 | Alex |
|---|
orders
| id | cust. |
|---|---|
| 3207 | 32 |
| 3209 | 32 |
分片 2customers
| 6 | Jeff |
|---|
orders
| id | cust. |
|---|---|
| 3203 | 32 |
| 3208 | 2 |
分片 3customers
| 32 | Maja |
|---|
orders
| id | cust. |
|---|---|
| 3202 | 2 |
| 32 | 1 |
分片 4customers
| 1 | Leah |
|---|
orders
| id | cust. |
|---|---|
| 3201 | 6 |
| 1 | 6 |
まず、4つの分片すべてから顧客を集める。
これは効率が悪い。もっと良い方法がある。
ハッシュ結合は別の結合手法で、今回の用途に向いている。考え方はこうだ。まず一方の入力からルックアップテーブルを構築し、それを使って、もう一方の入力の行が届くたびに一致するものを探す。同じように、ルーターはまずリクエストを1つ送って customers 行をまとめて取得するが、今回は顧客 ID をキーとするハッシュマップをルーターのメモリ上に構築する。そうすれば、9月9日以降に作成された orders 行をまとめて要求する大規模な scatter-gather を1回実行できる。
router
| id | cust. | total | date |
|---|---|---|---|
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| ··· | ··· | ··· | ··· |
| name | total | created_at |
|---|---|---|
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
| ··· | ··· | ··· |
SELECT id, nameFROM customers;
| key | name |
|---|---|
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
| ··· | ··· |
Shard 1customers
| 2 | Alex |
|---|
orders
| id | cust. |
|---|---|
| 3207 | 32 |
| 3209 | 32 |
Shard 2customers
| 6 | Jeff |
|---|
orders
| id | cust. |
|---|---|
| 3203 | 32 |
| 3208 | 2 |
Shard 3customers
| 32 | Maja |
|---|
orders
| id | cust. |
|---|---|
| 3202 | 2 |
| 32 | 1 |
Shard 4customers
| 1 | Leah |
|---|
orders
| id | cust. |
|---|---|
| 3201 | 6 |
| 1 | 6 |
ひとつの方法は、まず customers を一度取得し、その後は order が届くたびに、対応する customer を router で引くというものだ。
customer のハッシュテーブルをメモリに置けばリソースは食うが、注文が届いた時点ですぐ照合でき、customer ごとに問い合わせを投げ直さずに済む。join のルールはこの二つの方法の推定コストを比べる。どちらも照合は router の中で、Postgres が各 shard 上で返す行を使って行う。
planner は入力の順序を逆にすることも検討する。先に customers を取得するのは、このどちらの join にも必須ではない。
planner の選択肢はこれだけではない。merge join やバッチ nested-loop join なども使える。その役目は、手持ちの情報からクエリにとって最も効率のよい実行計画を見つけることにある。
このクエリでの選択
planner がどちらのテーブルを先に取得し、どの join 方式を使うかは、テーブルの大きさに大きく左右される。4 つの shard 全体で customers が 100,000 行、orders が 1,000,000 行あるとしよう。日付の不等式に対する既定の見積もりでは、およそ 3 分の 1 の注文が該当し、約 300,000 行になる。次の Picasso 図は、planner がどんな状況で異なる join アルゴリズムと順序を選ぶかを示している。行数を変えれば、planner の選択がどう変わるかが見える。Orders1101001,00010,000100,0001M10M 1101001,00010,000100,0001M10MCustomers__Nested loop · orders first__Nested loop · customers first__Hash join · build customers__Hash join · build orders__Batched nested loop この行数でなぜ hash join を選ぶのかを見てみよう。先ほどの nested loop 方式なら、100,000 回の customer 検索(すべてを検索する)を行い、そのうえで customer ごとに orders へ scatter-gather をかけることになる。取得する行はおよそ 400,000 行(customers 100,000 行と orders 300,000 行)だが、router と shard の間に何千もの個別クエリとして散らばる。hash join なら、customers を 4 回の shard クエリで、該当する orders をさらに 4 回で取得する。この 8 回のクエリでも約 400,000 行が router に渡り、その中には結局マッチしない customers も含まれる。この方法には手間がひとつ増える。メモリ上でハッシュテーブルを構築するのだ。planner はインデックス列の異なる値の個数も数え、等値フィルタや join が何行にマッチするかを見積もる。その見積もりをもとに、shard へのリクエスト数、転送する行数、処理する行数、メモリに常駐する行数を天秤にかける。この行数では、customer の入力を先に取ってローカルに参照用として保持するほうが、リクエストを繰り返し投げるより推定コストがはるかに低い。ここで planner にはもうひとつ選ぶことがある。ハッシュテーブルはどちらの入力からでも構築できる。planner は customers 100,000 行のメモリ所要量が、該当する orders 300,000 行より小さいと見積もるので、customers を選ぶ。ここまで決まれば、最初の計画を最終的なクエリ計画に変えられる。customers の経路がハッシュテーブルの構築を担い、orders の経路が探索に使う行を供給する。どちらも scatter 経路であり、つまりそれぞれがクエリを 4 つの shard すべてに送る。どちらにも、リクエストを特定の shard-key 値に絞り込む条件はない。ではハッシュテーブルがメモリに収まらなかったらどうなるか。メモリは速くて便利だが、使い切らなければの話だ。ハッシュテーブルがメモリの予算を超えると、router はディスクへ溢し出す。join key で customers と orders をパーティションに分け、対応するパーティション同士を順に join する。
*模式的な行データと小さめのメモリ予算を示した図で、表しているのはスピルであり、記事の例が必要とするメモリではない。シャードはルーターの下、アプリケーションはその上、一時ディスクストレージはルーターの RAM の右側にある。青い customer 行はシャードから上へ流れ、ルーターのハッシュテーブルを埋める。メモリ予算に達すると、3 つの一時ディスクパーティションへスピルする。黄色い orders はシャードから届き、RAM 内の書き込みバッファを通ってから、一致する結合キーごとにディスクパーティションへフラッシュされる。両方の入力のパーティション化が終わると、ルーターは customer パーティションを 1 つ RAM に読み込み、対応する orders でプローブし、一致を返してから次のパーティションへ進む。ディスク上のコピーは残り、メモリ内のハッシュテーブルは再利用される。模式的なメモリメーターは、青い customer メモリ、黄色い order の書き込み/読み取りバッファ、処理メモリを示す。orders は保持された customer ハッシュテーブルを通ってストリーム形式で流れる。order パーティション全体が読み込まれることはない。このメーターは実測バイト数でも、正確なメモリ予算の会計カウンターでもなく、ルーターの他の処理は含まない。メモリ内の各 customer パーティションは結合が完了するとクリアされ、それから次のパーティションが読み込まれる。黄色い order はそれぞれ customer_id、total、created_at を RAM とディスクを通じて運ぶ。結合結果には customers.name、orders.total、orders.created_at が含まれる。タイムスタンプはコンパクトな日付で表示される:September 9, 2026。例の時間範囲は 12:00 から 17:00 UTC。6 件の一致結果はすべてアプリケーションへ上って渡される。一時停止の後、ループが繰り返される。
ハッシュテーブルがメモリに収まってもディスクへスピルしても、結合はルーター内で起きる。これで計画ができた。実行しよう。
実行
まだルーターから離れていないが、まもなくクラスタ内の別のノードへ進む。これがすべて単一インスタンスで起きるなら、Postgres はローカルでテーブルを結合できる。しかし我々の計画は結合をルーターに置き、ルーターはどちらのテーブルも保持していない。_シャーディングのおかげで。_計画に必要な行は 4 つの独立したインスタンスに散らばっているので、ルーターは今、シャードへ個別の customer クエリと order クエリを送り、返ってきた行を結合しなければならない。まず、4 つのシャードそれぞれが customer クエリを実行する必要がある:
SELECT customers.id, customers.name
FROM public.customers;
ルーターが customer ハッシュテーブルを構築したら、各シャードが次のクエリの結果を実行して返す:
SELECT orders.customer_id, orders.total, orders.created_at
FROM public.orders
WHERE orders.created_at >= $1;
order クエリは 4 つのシャードすべてで同じ $1 タイムスタンプパラメータ '2026-09-09 00:00:00+00' を使う。最初の 4 つのリクエストはすでにルーターから送り出せる。
Postgres への接続
customer リクエストは 4 つのシャードすべてへ送られる。ここでは shard 2 で起きていることだけを見る。ルーターは Postgres に直接接続せず、リクエストを sidecar へ送る。sidecar は Postgres と並行して動く別プロセスだ。これで Postgres との間に 1 ステップ増えるが、接続管理を各インスタンスのそばに集約でき、異なるルーターからのリクエストが同じ Postgres 接続プールを共有できる。sidecar はこれらの接続を安全に再利用できることを保証する。もし Postgres へ直接接続していたら、クエリはあらかじめ確立されたセッションで実行される。しかし我々のアーキテクチャでは、そのセッションはルーターに属する。シャードへのリクエストは、データベース、認証ロール、セッション設定をクエリと一緒に運ぶ必要があり、そうして各部分がアプリケーションの期待する権限と設定で動く。ルーターは拡張 Postgres プロトコルメッセージの列を用意する:Parse、Bind、Execute、Sync。今回は Parse に最初の join ではなく customer クエリが入る。ルーターはこれらのメッセージを ExecuteRequest の raw フィールドに詰め、宛先アドレスとセッション情報も一緒に持たせる。リクエストは protobuf を使い、Postgres メッセージをバイトのまま運ぶ。
次にリクエストを送るチャネルが要る。ルーターは sidecar ごとに、長寿命の双方向 gRPC ストリームのプールを維持する。空きストリームを 1 つ借り、なければ新しく開き、ExecuteRequest を送る。sidecar は同じストリームでレスポンスチャンクを返す。ルーター
生の Postgres メッセージ
-
Parse
SELECT customers. id, customers.name
FROM public .customers
-
Bind
parameters none 3. #### Execute
max_rows = 0 4. #### Sync
targetshard 2 · primary
databaseapplication database
session_userauthenticated role
options.settingssession settings
SidecarShard 2 · primary
クライアントリクエストはようやく shard 2 に届いたが、まだ Postgres には入っていない。
Sidecar の内部
我々は shard 2 の sidecar まで来たが、まだ Postgres 内で動いてはいない。Postgres の接続ごとに 1 プロセスというモデルがまた立ちはだかる。忘れてはならないが、各 Postgres 接続には独立したサーバーバックエンドプロセスが必要で、アイドル状態でもメモリを占める。各 router が各 shard への接続プールをそれぞれ維持するとしたら、規模が大きくなるにつれてこれらのプロセスは急速に積み上がる。sidecar はこれらの router に、必要に応じて借りられるホットな Postgres 接続プールを共有させる。我々のクライアントクエリはそのうちの 1 接続を使えるが、その接続には前回の使用が残した設定とアイデンティティがまだ付いているかもしれない。ここで、SQL と一緒に送られるセッション情報が役に立つ。Router は shard 2 の sidecar への利用可能な gRPC stream を 1 本借りる。我々のリクエストは SQL、Settings、アイデンティティを運び、sidecar のインターフェース上で一緒に表示される。既存の Postgres 接続が Sidecar の独立したプールから立ち上がる。その以前の設定とアイデンティティは、リクエスト内の状態に道を譲る。まず設定を揃え、次にアイデンティティを揃える。両方が整ってから、SQL が Postgres、すなわち shard 2 の primary へ送られる。すでに一致している状態はそのまま再利用でき、変更は不要。Stream と Postgres 接続は恒久的に結び付いているわけではない。各プールにはそれぞれ並行するチャネルがある。プールの占有状況とタイミングは模式的なものにすぎない。立ち上がりとドッキングは、接続を取り出してセッションを準備する過程を表しており、物理的な接続を新規作成するのでも、接続プール自体を移すのでもない。この説明はループで繰り返され、レスポンスや解放は表さない。
SELECT customers.id, customers.name FROM public.customers
設定セッション設定ID認証済みロール
設定前回の使用セッション設定ID前回の使用認証済みロール
コネクションプールはまず、リクエストに合う設定の空きコネクションを探す。見つかればそのまま使える。見つからなければ、sidecar が以前の設定をデフォルトに戻し、こちらの設定を適用する。
sidecar はセッションのロールも確認し、正しいことを検証する。ロール汚染を防ぐためだ。直前にそのコネクションを借りていたのが誰であれ、こちらのクエリはこちらの権限で実行されなければならない。
sidecar が SQL を転送し、ようやくクライアントのクエリが Postgres に届く。
ようやく Postgres に到達
クライアントのクエリが Postgres に届いた。
SELECT customers.id, customers.name
FROM public.customers;
router の層でこれだけ計画に時間を費やしてきたが、ここで Postgres も独自の計画を立てる。本物の Postgres インスタンスなので、自前の planner と統計情報を使って customers テーブルの読み取り方法を決め、shard 2 上の customers に対してその計画を実行する。行の準備ができると、Postgres は借りたコネクションを通じて sidecar に送り返す。ハッシュテーブルの原型ができた。
router に戻る
sidecar は、送信リクエストで使った gRPC ストリームを通じて行を router に中継する。行はバッチで届くので、shard がまだデータを返している間に router は customer ハッシュテーブルの構築を始められる。4 つの customer リクエストがすべて完了し、ハッシュテーブルができあがると、ハッシュ結合が orders Route の計画ノードを起動する。さらに 4 つのリクエストが、今たどってきたばかりの転送とコネクションプールの経路を進む。今度は Postgres が各 shard で 9 月 9 日というフィルタ条件を適用し、該当する注文を返す。返ってきた行でマッチが発生した。router は各注文の customer_id を使ってハッシュテーブルから対応する customer を見つけ、顧客名、注文合計、作成タイムスタンプを含む結果行を組み立てる。ここでかなり巧妙な工夫があり、計算サイクルをいくらか節約できる。Postgres に出力フィールドを要求するとき、クライアントが元のクエリで選んだのと同じテキスト形式またはバイナリ形式を使うのだ。値はすでにエンコードされている。router はマッチした行からエンコード済みの名前、合計、タイムスタンプをメモリ上で直接コピーして結果を組み立てる。Go の値にデコードしてから再エンコードする手間が省ける。customersid 32 · Maja440000001A0002000000080000000000000020000000044D616A61
orderscustomer_id 32440000003500030000000800000000000000200000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
DataRow → driver_7 新ヘッダバイト_ + 43 コピー済み44000000310003000000044D616A610000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
000000044D616A610000000533302E303000000016323032362D30392D30392031323A30303A30302B3030
customer 行には ID 32 と Maja が含まれる。order 行には customer ID 32、合計 30.00、タイムスタンプ 2026-09-09 12:00:00+00 が含まれる。2 つのバイナリ ID はマッチングのために読み取られ、出力には入らない。エンコード済みの name、total、timestamp のコピーが、それぞれのフィールド長プレフィックスとともに、新しいメッセージ内の正確な位置へ移動する。これらの出力値はテキスト形式を使う。新しいヘッダは 44 00 00 00 31 00 03。長さは 49 で、D タグを含まず、列数は 3。完全な DataRow は 50 バイト。ソースのバイト列は変わらない。
各 order リクエストの完了を待つ必要もない。アプリの Postgres クライアントは、shard がまだ order を返しているうちに結果を受け取り始められる。
4 つの shard すべてが order を返し終えると、このクエリの旅は終わる。アプリは結果を受け取り、コネクションは次のクエリを待つ。
評価エンジン
このバイト直接コピーの近道が使えない場合もある。ここまで追ってきたクエリでは、出力のどの値も shard が返した行にすでに存在していた。router はエンコード済みのバイトをそのまま結果にコピーできる。だが、代わりに customer ごとの平均 order total を要求したらどうなるだろうか。
SELECT customers.name, AVG(orders.total) AS average_total
FROM customers
JOIN orders ON orders.customer_id = customers.id
WHERE orders.created_at >= $1
GROUP BY customers.id, customers.name;
これには router 内のあるコンポーネントが追加の作業をする必要がある。おそらく見覚えがあるだろう、Postgres の内部コンポーネント、Eval Engine だ。1 人の customer の order は 4 つの shard すべてに分散していることを思い出してほしい。1 つの shard だけでは完全な平均を計算できない。shard の平均をさらに平均することもできない。ある shard にはある customer の order が 2 件、別の shard には 20 件あるかもしれない。そこで router は平均を合計とカウントに書き換え、各 shard がローカルに計算する。結果が返ると router が customer ごとにマージする。router の評価エンジンがマージ後の合計をマージ後のカウントで割り、その customer の平均を出す。
ルーターは AVG(orders.total) を SUM(total) と COUNT(total) に書き換え、その処理を 4 つのシャードに振り分ける。各シャードはサンプルの顧客グループごとに合計と件数を計算する。顧客 32、つまり Maja について、シャード 1 は合計 $100・件数 2、シャード 2 は $90・3、シャード 3 は $80・1、シャード 4 は $50・2 を返す。件数には NULL でない注文合計だけが含まれる。ルーターはこれらを統合して $320・8 とする。評価エンジンは $320 を 8 で割り、平均 $40 を Maja の名前とともにアプリへ返す。4 つのシャードの平均をさらに平均すると、誤って $46.25 になる。これらのパネルが区別しているのは概念上の手順であり、物理プランの演算子ではない。この特写では顧客 join とその他のグループ化を省略している。各注文は 1 人の顧客にしか一致しない。数値・到着順・タイミングはすべて例示用。視差効果を減らす設定が有効な場合、計算完了後の結果は表示されたままになる。
とはいえ、最初のクエリに関して言えば、評価エンジンにそれほど仕事はない。
より良いトポロジー
ルーターがシャードをまたぐ join と計算にどれだけの処理を要するかは、すでに見たとおりだ。選んだシャードキーのせいで、その処理は必要以上に重くなっている。主キーでシャードすること自体が問題なのではない。顧客の分散には引き続き customers.id を使えばいい。複雑さの原因は、注文を orders.id で分散していることにある。これは必要な顧客とは独立していて、join の際に対応づけなければならない。orders を customer_id でシャードし、customers と同じハッシュとシャードマッピングを使えば、各顧客は自分の注文と同じ場所に留まる。Postgres は join をローカルで完結できる。クエリは全顧客の直近注文を求めるため、依然として 4 つのシャードすべてに届く。だが今回は完全な join を各 Postgres インスタンスへ送り、それぞれが自分の行をローカルで join して結果を返す。シャードへのリクエストは 8 回ではなく 4 回で、ルーターが顧客のハッシュテーブルを構築する必要もなくなる。顧客は引き続き customers.id でシャードされる。ここで使っているのは記事の前半と同じサンプル行だ。注文は customer_id を使い、customers.id と同じシャードマッピングを適用する。アプリは同じクエリをルーターへ送る。ルーターは join 済みのクエリを 4 つのシャードそれぞれへ送る。Postgres は一致する顧客と条件を満たす注文をローカルで join する。シャード 4 もクエリを受けるが、条件を満たす注文がないため join 後の行は返さない。他のシャードは合計 6 行の join 結果を返し、ルーターは届いた順にそのまま転送する。どちらのレイアウトでも、ルーターからアプリへ返る結果は同じ 6 件で、その後は完了した結果の上で一時停止し、再生を繰り返す。このサンプル行の移動は 2 つのレイアウトを比較するためのもので、実際のリシャーディングの過程ではない。アニメーションの速度はベンチマークではない。
- シャード 1 は顧客 2 と、顧客 2 用に選ばれた注文 3202、3208 を保持する。
- シャード 2 は顧客 6 と、顧客 6 用に選ばれた注文 3201 を保持する。
- シャード 3 は顧客 32 と、顧客 32 用に選ばれた注文 3207、3203、3209 を保持する。
- シャード 4 は顧客 1 を保持し、選ばれた注文はない。
先ほど説明した認証、転送、接続管理は依然として必要だ。だが、別々に取得した顧客と注文を router 内で突き合わせる作業は消える。最速の router join とは、そもそも実行しなくていい join のことだ。これは数千語も前に言っておくべきだった。
ほら、簡単だろう?
Neki へようこそ。