はじめに
このガイドでは、まず ClickHouse が分散テーブルを介してクエリを 複数の分片にどのように分散するのかを説明し、その後、クエリの実行で 複数のレプリカをどのように活用できるのかを説明します。
分片アーキテクチャ
shared-nothing アーキテクチャでは、クラスターは通常、複数の分片に 分割され、各分片には全体データの一部が格納されます。これらの分片の 上位には分散テーブルがあり、データ全体を単一のビューとして提供します。 読み取りはローカルテーブルに送ることができます。その場合、クエリの実行は 指定した分片でのみ行われます。あるいは分散テーブルに送ることもでき、その 場合は各分片で指定されたクエリが実行されます。分散テーブルに対するクエリを 受けたサーバーがデータを集約し、クライアントに応答を返します。 上の図は、クライアントが分散テーブルにクエリを実行した際に何が起こるかを示しています。- SELECT クエリは、任意のノード上の分散テーブルに送信されます (ラウンドロビン戦略を介するか、ロードバランサーによって特定のサーバー にルーティングされた後)。このノードがコーディネーターとして動作します。
- このノードは、分散テーブルで指定された情報に基づいて、 クエリを実行する必要がある各分片を特定し、クエリを 各分片に送信します。
- 各分片はローカルでデータを読み取り、フィルタリングと集約を行った後、 マージ可能な状態をコーディネーターに返します。
- コーディネーターとなるノードがデータをマージし、その後 応答をクライアントに返します。
非分片アーキテクチャ
並列レプリカの紹介
複数のサーバーにまたがってクエリ実行を並列化するには、まず サーバーのうち 1 台をコーディネーターとして割り当てられるようにする 必要があります。コーディネーターは、実行すべき タスクの一覧を作成し、それらがすべて実行・ 集計され、結果がクライアントに返されることを保証する役割を担います。多くの 分散システムと同様に、通常これは最初の クエリを受け取ったノードが担います。また、作業単位を定義する必要もあります。分片アーキテクチャでは、 作業単位は分片、つまりデータの一部分です。並列レプリカでは、 作業単位として グラニュール と呼ばれるテーブルの小さな部分を 使用します。 では、以下の図を使って、実際にどのように動作するのか見てみましょう。 並列レプリカでは:- クライアントからのクエリは、ロード バランサーを経由して 1 つのノードに送られます。このノードがこのクエリのコーディネーターになります。
- そのノードは各パーツの索引を解析し、処理対象のパーツと グラニュールを選択します。
- コーディネーターは、ワークロードを 異なるレプリカに割り当て可能な一連のグラニュールに分割します。
- 各グラニュールの集合は対応するレプリカで処理され、完了すると マージ可能な状態がコーディネーターに送られます。
- 最後に、コーディネーターがレプリカからのすべての結果をマージし、 クライアントにレスポンスを返します。
- 利用できないレプリカがある場合があります。
- ClickHouse のレプリケーションは非同期であるため、ある時点では一部のレプリカが 同じパーツを持っていない可能性があります。
- レプリカ間のテールレイテンシを何らかの方法で扱う必要があります。
- ファイルシステムキャッシュは各レプリカ上の アクティビティに応じて異なるため、タスクをランダムに割り当てると、 キャッシュ局所性の観点では最適でないパフォーマンスにつながる可能性があります。
アナウンス
上のリストの (1) と (2) に対処するため、アナウンス という概念を導入しました。以下の図でその仕組みを見てみましょう。- クライアントからのクエリは、ロード バランサーを経由して 1 つのノードに送られます。このノードが、このクエリのコーディネーターになります。
- コーディネーターとなるノードは、 クラスター内のすべてのレプリカからアナウンスを取得するためのリクエストを送信します。レプリカごとに、テーブルの現在のパーツ集合の見え方が わずかに異なる場合があります。そのため、 不正確なスケジューリング判断を避けるには、この情報を収集する必要があります。
- その後、コーディネーターとなるノードはアナウンスを使って、 各レプリカに割り当て可能な グラニュール の集合を定義します。たとえばここでは、 パーツ 3 の グラニュール はレプリカ 2 には 1 つも割り当てられていないことがわかります。 これは、このレプリカがアナウンスでこのパーツを通知しなかったためです。 また、 レプリカ 3 はアナウンスを返さなかったため、タスクが 1 つも割り当てられていない点にも注意してください。
- 各レプリカが自分に割り当てられた グラニュール のサブセットに対してクエリを処理し、 マージ可能な状態がコーディネーターに送り返されると、 コーディネーターが結果をマージし、そのレスポンスがクライアントに送信されます。
動的協調
テールレイテンシの問題に対処するため、動的協調を追加しました。これは、 すべてのグラニュールを1回のリクエストでレプリカに送るのではなく、各レプリカが コーディネーターに新しいタスク (処理するグラニュールのセット) を要求 できることを意味します。コーディネーターは、受け取ったアナウンスに基づいて グラニュールのセットをレプリカに割り当てます。 ここでは、すべてのレプリカがすべてのパーツを含む アナウンスを送信し終えた段階にあると仮定します。 以下の図は、動的協調がどのように機能するかを示しています。- レプリカは、タスクを処理できることをコーディネーターノードに知らせます。 また、どれだけの作業を処理できるかを指定することもできます。
- コーディネーターはレプリカにタスクを割り当てます。
- レプリカ1と2は、自分たちのタスクを非常にすばやく完了します。 そのため、コーディネーターノードに別のタスクを要求します。
- コーディネーターはレプリカ1と2に新しいタスクを割り当てます。
- すべてのレプリカが、それぞれのタスクの処理を完了しました。 さらにタスクを要求します。
- コーディネーターは、アナウンスを使って未処理のタスクが残っているかどうかを 確認しますが、残っているタスクはありません。
- コーディネーターは、すべての処理が完了したことをレプリカに伝えます。 その後、マージ可能なすべての状態をマージして、クエリに応答します。
cache の局所性の管理
| レプリカ 1 | レプリカ 2 | レプリカ 3 | |
|---|---|---|---|
| パート 1 | g1, g6, g7 | g2, g4, g5 | g3 |
| パート 2 | g1 | g2, g4, g5 | g3 |
| パート 3 | g1, g6 | g2, g4, g5 | g3 |
max_parallel_replicas がレプリカ数より少ない場合は、クエリ実行のためにランダムにレプリカが選択されます。
タスクのスティーリング
制限事項
この機能には既知の制限があり、主なものをこのセクションで説明します。以下に挙げた制限事項以外の問題が見つかり、その原因が並列レプリカにあると思われる場合は、
GitHub でラベル
comp-parallel-replicas を付けて issue を作成してください。並列レプリカの問題の調査
各クエリでどの設定が使われているかは、system.query_log テーブルで確認できます。また、
system.events
テーブルを見ると、サーバー上で発生したすべてのイベントを確認できます。さらに、
clusterAllReplicas テーブル関数を使うと、すべてのレプリカ上のテーブルを確認できます
(Cloud ユーザーの場合は default を使用してください) 。
Query
レスポンス
レスポンス
Response
system.text_log テーブルには、並列レプリカを使用したクエリの実行に関する情報も
含まれています。
Query
レスポンス
レスポンス
Response
EXPLAIN PIPELINE を使用することもできます。これにより、ClickHouse
がクエリをどのように実行し、クエリの実行に
どのようなリソースが使われるかを確認できます。例として、次のクエリを見てみましょう。
EXPLAIN PIPELINE (without parallel replica)
EXPLAIN PIPELINE (with parallel replica)