# D3S: Debugging Deployed Distributed Systems > [!abstract] 概要 > 大規模分散システムのテストは難しい。一部の誤りは、マシン障害やネットワーク障害を伴う分散した事象の並びの後にだけ現れるからである。D3S は、開発者が本番稼働中のシステムの分散的な性質に関する述語を指定でき、システムの稼働中にその述語を検査するチェッカーである。D3S が問題を見つけると、その問題に至った状態変化の列を出力するため、開発者は根本原因をすばやく特定できる。開発者は述語を単純な逐次的プログラミングのスタイルで書き、D3S はその述語を分散・並列に検査する。これにより検査は大規模システムに対してスケールし、耐障害性も備える。バイナリ計装を用いることで、D3S はレガシーシステムに対して透過的に動作し、検査する述語を実行時に変更できる。5 つの本番系システムでの評価により、D3S は自明でない正しさの誤りと性能の誤りを実行時に低い性能オーバーヘッド(8% 未満)で検知できることが示された。 ## 論文情報 - 著者: [[Xuezheng Liu]]、[[Zhenyu Guo]]、Xi Wang(清華大学)、Feibo Chen(復旦大学)、Xiaochen Lian(上海交通大学)、Jian Tang、Ming Wu、[[M. Frans Kaashoek]]([[MIT CSAIL]])、Zheng Zhang(Xuezheng Liu・Zhenyu Guo・Jian Tang・Ming Wu・Zheng Zhang は [[Microsoft Research Asia]])。所属は論文の脚注記号に基づく。 - 掲載: NSDI 2008(5th USENIX Symposium on Networked Systems Design and Implementation)、pp. 423-437。シェパードは Amin Vahdat。 - 実装対象は Windows。バイナリ計装は WiDS BOX(Detours 系のツールキット)を用いる。 ## 概要 [[D3S]] は、本番稼働中の分散システムに対して、開発者が書いた述語(分散不変条件)を実行時に検査するチェッカーである。述語は逐次的な C++ 関数として書き、実行系が(1)バイナリ計装による状態の露出、(2)論理時計による大域スナップショットの構築、(3)キー空間の分割による並列検査、(4)検査側の障害からの再実行を担う。違反時は違反に至った状態変化の列を返すため、事後のログ解析より根本原因へ速く届く。PacificA、Paxos 実装、検索エンジン、Chord、BitTorrent クライアントの 5 系統で、非自明なバグと性能問題を見つけた。 ## 問題設定 - 大規模分散システムのバグは、マシン障害やネットワーク障害を含む特定の事象列の後にだけ現れ、再現が難しい。 - 実務の手法は、開発者が print 文で局所状態を吐き、バッファを中央に送り、スクリプトで各マシンの状態を並べて大域スナップショットにして検査する方式である。この方式には次の欠点がある。 - 状態の記録と大域的に整合したスナップショットへの並べ替えを、開発者が自前で書く。 - どの状態を記録するかを事前に予測する必要がある。多く記録すれば本番系が遅くなり、少なければ誤動作を見逃す。 - 中央のチェッカーは追いつかないことがある。著者らが扱ったあるアプリケーションは、マシンあたり 500〜1000 KB/s の監視データを出した(全データの 1〜2% だが、全体では 1 台で処理できない量である)。 - チェック対象のプロセスが落ちたときに、大域スナップショットをどう近似するかと、チェック自体の耐障害化を考える必要がある。 - 課題は 3 つに整理される。(1)何を検査するかを簡単に書けたまま、複数マシンで検査できること。(2)検査マシンの障害を扱うこと。(3)検査対象プロセスの障害を扱い、偽陽性・偽陰性を不必要に出さないこと。たとえば、リース付きでロックを取ったクライアントがロック解放前に落ち、リース失効後に別のクライアントがロックを取った場合を、二重取得として誤報してはならない。 ## 提案手法 図 1 に全体像を示す。述語をコンパイルして、状態露出ライブラリ(state exposer)と検査ライブラリの 2 つの DLL を作る。状態露出ライブラリをバイナリ計装で稼働中のプロセスへ注入し、検査ライブラリを検証プロセス(verifier)へ注入する。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig01-overview.png]] *Figure 1: Overview of D3S.* ### 述語の書き方(§2.1、2.3) 述語は、頂点が計算段階で辺がデータの流れである非巡回グラフ(Dryad に着想を得た構造)として組む。頂点 V0 は検査対象システム自身が出す状態タプル(例: クライアント ID、ロック ID、モード)、後続の頂点は逐次的な C++ の `Execute` メソッドで書く。ロック整合性の例では、`Execute` が大域スナップショット中のタプルを数え上げ、あるロックに排他保持者が 2 以上、または排他保持者と共有保持者が併存する場合を衝突として出力する。型定義は検査対象システムのヘッダーを再利用できる。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig02-lock-predicate.png]] *Figure 2: (a) Checking code. (b) graph and checker execution.* ### 述語の挿入(§2.2) コンパイラが状態露出用と検査用の 2 つの DLL を生成する。状態露出側は、対象関数の前後に計装コードを差し込み、関数の引数(`$0` は `this`)からタプルを作り、`addtuple` / `deltuple` で V0 の出力に加除する。著者らの経験では、検査した全システムで関数引数の露出だけで状態変化の観測に足りた。稼働中に述語を追加すると、それ以前の未監視の履歴に依存する違反は見逃しうる。 ### 分割実行とデータフロー(§2.3、2.4) 各時刻 t の処理は、上流からの入力がそろって整合的なスナップショットが組めた時点で走り、パイプライン的に並列化される。述語は露出状態から決定的に計算されるため、途中頂点の障害後は V0 のバッファ済み状態から同じ時刻を再実行できる。並列化のため、`Mapping` メソッドで出力タプルを仮想キー空間へ写像する(MapReduce の Map に相当)。同じロックのタプルを同じ検証プロセスへ集める、といった要件をここで表す。キー範囲の割当は動的に変更でき、検証プロセスの追加・削除・負荷再配分・障害時の引き継ぎに使う。 ### ストリーム処理と標本化(§2.5) 隣接時刻の状態差は小さいため、頂点は直前時刻との差分だけを送り、任意の `ExecuteChange` で増分計算する。Chord の鍵範囲述語では、V0(各ノードの pred / self / succ)、V1(範囲の総和を部分集合ごとに計算)、V2(分割ごとの和を集約)の 3 段の集約木を使う。加えて、キー空間の一部だけを検査する標本化と、時刻の標本化ができる。これらは軽量化の代わりに偽陰性の恐れを受け入れる。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig03-chord-range-coverage.png]] *Figure 3: The predicate that checks the key range coverage among Chord nodes.* ### 大域スナップショットと述語の正しさ(§3) - 実行を、大域タイムスタンプ付きの整合スナップショットの列 π(t) をたどる逐次状態機械とみなす。π(t) は時刻 t のメンバーシップ M(t) の各プロセスの局所状態 Sp(t) の和集合である。 - 大域タイムスタンプには [[@1978__CACM__Time Clocks and the Ordering of Events in a Distributed System|Lamport の論理時計]]を使う。時計読み取りごとに 1 増やし、メッセージに時計値を添え、受信側は自分の値と添付値の最大を取る。これで happens-before を保ち、全タプルを整合的な全順序に並べる。 - 述語は連続する n 個のスナップショット(窓幅 n)に対する関数 P(ti) = F(π(ti−n+1), ..., π(ti)) と定義する。著者らの経験では、有用な性質はすべて直近の窓だけで検査できた。 - 検証プロセスは、障害検知器に M'(t) を問い合わせ、M'(t) の全プロセスの Sp(t) がそろうまで待って π(t) を作る。あるプロセスが長時間状態を出さない場合に備え、状態露出側が現在のタイムスタンプを定期的にハートビートとして送る。ハートビートは障害検知器と進捗通知を兼ねる。障害と判定されたプロセスの状態は、判定時刻以降のスナップショットから除く。 - 障害検知器の出力が正しければ(M'(t) = M(t))、スナップショットは完全である。不完全なスナップショットは偽陽性・偽陰性を生む。バッファ時間 Tbuf は、警報の速さと正確さのトレードオフを決める調整項で、検知器の Tout より長くする。PacificA では、最大メッセージ遅延 350ms、キープアライブ 1000ms から Tbuf を 2000ms とした。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig04-snapshot-membership.png]] *Figure 4: Predicate checking for consistency of distributed locks.* ### 実装(§4) - WiDS BOX が関数アドレスをシンボル情報から解決し、コードモジュールの関数入口を書き換えて呼び出しを DLL のコールバックへ向ける。書き換えは原子的で、計装中にスレッドを停止しない。 - 露出状態は内部バッファに溜め、500 バイトを超えるか、最後の送信から 500ms 経過したら送る。ソケット API も横取りして論理時計を更新し、メッセージごとに 8 バイトを付加する。論理時計は 1 秒を 1000 ティックとし、実時間と対応づける。 - 状態露出側と検証プロセスの間、および検証プロセス間は信頼できる伝送を用いる。中央のマスターがキー空間の分割を管理し、検証プロセスが報告する最新検証時刻が途絶えると障害とみなして分割を張り替え、新しい分割を全体へ配る。 ## 新規性 - 述語を逐次コードで書く単純なモデルと、それを大域スナップショット上で分散・並列に走らせる実行系を一体にした点。従来の再生型検査や P2 monitor は、状態の収集と順序付けを利用者に任せていた。 - 障害プロセスをスナップショットから除く形式化(メンバーシップと障害検知器)で、検査対象の障害時にも誤報を抑える点。 - バイナリ計装により、レガシーで本番稼働中のシステムを無改造で検査でき、述語を稼働中に差し替えられる点。検査側の障害は、決定的な検査を同じ入力で再実行して隠す。 - 関連研究との違い: 再生型検査(WiDS checker、Friday)は全実行の再生が大規模では高価であり、D3S は再生を補完する位置づけとする(オンライン検査で違反箇所を絞り、後で限定的な再生を行う展望)。P2 monitor は OverLog 製のシステムに限られる。モデル検査は状態爆発で小規模にとどまり、環境が仮想的で性能バグを見つけにくい。ログ解析は事後であり、D3S は「状態を拡張可能に集める賢いロガー」に加えて実行時に露出状態を変えられる点が利点とする。並列アプリ向けのスタックトレース比較などは同質のプロセス向けであり、D3S は異質なプロセスにも使える。 ## 実験設定 - 評価は、本番品質のコードを含む 5 つの稼働中システムで、1 か月以内に得た結果である。Table 1 に対象を示す(LoC はシステムのコード行数、LoP は述語のコード行数)。 - 既定の環境は、2 GHz Xeon 2 基・4 GB メモリ・1 GbE・Windows Server 2003 のマシンである。一部のシステムでは網羅性のため障害を注入した。 *Table 1: Benchmarks and results* | 対象 | LoC | LoP | 述語 | 結果 | |---|---|---|---|---| | PacificA(構造化ストレージ) | 67,263 | 118 | メンバーシップ群の整合性、レプリカの整合性 | 正しさのバグ 3 件 | | MPS(Paxos 実装) | 6,993 | 50 | 合意出力の整合性、リーダー選出 | 正しさのバグ 2 件 | | Web 検索エンジン | 26,036 | 81 | インデックスサーバの応答時間の不均衡 | 性能問題 1 件 | | i3-Chord(DHT) | 7,640 | 72 | 鍵範囲の総被覆、鍵保持者の衝突 | 可用性と整合性の測定 | | libtorrent(BitTorrent クライアント) | 36,117 | 210 | 近傍集合、ダウンロード済みピース、ピア貢献度順位 | 性能バグ 2 件、フリーライダー検出 | - オーバーヘッドの評価は、状態露出の頻度と大きさを変えた、CPU 飽和状態のマイクロベンチマークと、PacificA の負荷実験である。 ## 実験結果 ### PacificA(§5.1) - スライス(100 MB)を 3 つの SliceServer に複製し、1 つを primary、残りを secondary とする。primary 不変条件は「任意の時刻でスライスあたり primary は高々 1」である。タプル `(Sid, MachineID, {P/S})` を露出し、Sid に写像して検査した。対象は SNS の社会グラフ(38 GB、途中表は 1 TB 超)を扱う 8 ストレージノード構成で、3 台の検証機を使った。 - 通常時は違反ゼロだった。SliceServer と MetaServer に障害を注入して数十回試すと、2 つの primary を持つ違反が検出された。状態列を追うと、MetaServer が要求を受理して primary を昇格させた後に落ち、復旧後に前回の応答を忘れて 2 つ目のレプリカの要求も受理したことが原因だった。受理した要求のログをバックグラウンドスレッドでバッチ書き込みしており、応答前にログを flush していなかった。 - 分散性質なので局所検査では捕まらない。事後ログ検証と違い、常時の部品レベル仕様検査で希少なバグを根本原因つきで捕まえた。古い版の 2 件のデータ競合は、開発者が数日かけて直したのに対し、D3S は数時間の通常利用で捕まえ、状態列が原因の手がかりになった。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig05-pacifica.png]] *Figure 5: PacificA architecture and the bug we found.* ### MPS Paxos(§5.2) - 5 ノードに配備し、各ノードが独立に API を呼ぶ。プロセスを落とす障害モジュールを状態露出と一緒に注入した。合意の安全性と活性(全ノードが同一の命令列を、連番が 1 ずつ増えて実行する)の述語は、数日間違反を検出しなかったが、稀な経路のポインタ誤りは障害注入で発見した。 - 仕様が緩いリーダー選出について、ノードの状態(Leader / Follower / Promoting / Learning)が正常状態から逸脱し続けることを検査する述語を書いた。ランダム障害注入下で数時間走らせると、Promoting のまま長時間停滞する C を検知した。A・B・E が落ちて復旧した後に C と D が同時に昇格を試み、D が過半数の承認を得る一方、C への承認メッセージの一部が失われる。B は D に対して遅れているため Learning に移り、C は A・E を死亡と考えたまま B にだけ要求を出し続け、Learning の B は状態転送以外を無視するため C が停滞する。外部の振る舞いは正しいが、Paxos の性質を直ちには壊さず、系を将来の障害に脆弱にする。開発チームが確認し、ログ検証では見逃されていた。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig06-paxos-leader-election.png]] *Figure 6: Leader election in Paxos implementation and the snapshots of states that lead to the bug.* ### Web 検索エンジン(§5.3) - 250 サーバ構成のバーティカル検索。フロントエンドが問合せを配り、インデックスサーバは文書人気で 5 階層に分けて上位層から検索する。検証機は 1 台。クエリの重要経路上の関数に実行時間を露出した。 - 負荷分散の述語 Pr = (max − mean) / mean を階層ごとに計算した。10 台のインデックスサーバ、500 問合せで、Pr が 1 を超える問合せ(max > 2 × mean)では、常に 1〜2 台が他より大きく遅かった。他の述語でさらに掘ると、単純なハッシュではインデックスサーバ間の負荷が均等にならず、複数語の問合せで性能を大きく落とすとわかった。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig07-search-imbalance.png]] *Figure 7: Load imbalance in the web search engine.* ### Chord(§5.4) - i3 の Chord 実装(Windows へ移植)を、7 台のマシン上の 85 ノードで実行した。図 3 の鍵範囲総被覆(100% 未満は穴で不可用、100% 超は重複で不整合)を、第 1 段 4 台・最終段 1 台の検証機で測った。近傍数を 3 と 8 の 2 構成にし、70% のノードを落として推移を見た。 - 3 近傍の構成は整合性と可用性の両面で挙動が不安定で、総被覆は 100% 前後で振動してから収束した。8 近傍は 100% を超えなかった。振動は、前任・後継をすべて失ったノードがフィンガーテーブル経由で再参加し、鍵範囲の重複を生むことによる。 - 総被覆は穴と重複が打ち消し合いうるため、鍵空間から 256 点を標本化し、保持ノード数を数える述語を足した。8 近傍でも総被覆が 100% を超えないのに、重複(鍵番号 200 付近)が起きていた。可用性と整合性のトレードオフを、集約木と標本化で実時間かつ定量的に測れる。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig08-chord-results.png]] *Figure 8: i3-Chord experimental results.(Figure 8a: aggregate key ranges. Figure 8b: holders of random sampled keys.)* ### BitTorrent(§5.5) - libtorrent 0.12 を 8 台・57 ピアで実行し、52 MB のファイルを配布した(アップロード帯域上限 10〜200 KB/s、完了に 15 分〜2 時間)。 - 近傍リストを集約するピアグラフの述語で、近傍が 300 を超えるピアを見つけた(ピアは 57 のみ)。原因は、近傍リストに IP を加える際の重複確認の欠落による重複 IP で、選択肢が減り性能が落ちていた。 - ピース分布を集約すると、一部のピースだけが全ピアへ速く広がり、ランダム選択の想定と食い違った。実装は、レプリカ数が同じピースについて全クライアントが同じ順序で選ぶ決定的な選択をしていた。修正後は、ピースごとの進捗がランダム選択に近づいた。 - EigenTrust を述語で実装し、貢献度を集中計算してフリーライダーを検出した。アップロード帯域を 20 KB/s に制限したピア 46〜56 を、他ピアと区別できた。デバッグ用途からオンライン監視へ用途を広げる例である。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig09-piece-distribution.png]] *Figure 9: Piece distributions over peers when finishing 30% downloading.* ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig10-peer-contribution.png]] *Figure 10: The contributions of peers (free riders are 46-56).* ### 性能(§6) - マイクロベンチマークでは、状態露出のオーバーヘッドは概ね 2% で、最大は 2 スレッドの場合(2 コアに一致し、追加の露出スレッドがスケジューリングを増やす)で 8% 未満である。Chord と Paxos のように CPU・I/O 負荷の低いシステムでは無視できる。BitTorrent と検索エンジンは 2% 未満である。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig11-overhead.png]] *Figure 11: Performance overhead on system being checked.* - PacificA は負荷に応じて変わるが 8% 未満である。1 台あたり平均 1,500 スナップショット/秒、ピーク時の追加帯域は 1000 KB/s 未満、状態露出は全 I/O の平均 0.5% 未満である。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig12-pacifica-overhead.png]] *Figure 12: Performance overhead on PacificA with different throughput.* - 60 秒時点で PacificA の全述語を新規に開始しても、クライアントのスループットに目立つ影響はない。検証機 3 台のうち 1 台を 30 秒時点で落とすと、キー範囲が再分割され、残りの検証機が負荷を引き継ぐ。 ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig13-throughput.png]] *Figure 13: Throughput of PacificA when a predicate starts.* ![[wiki/sources/_attachments/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems/fig14-verifier-failure.png]] *Figure 14: Load when verifier 1 fails at 30th second.* ## 考察 - 仕様や不変条件がある系(部品レベルの仕様がある良く設計されたシステム)では述語が書きやすく効果的である。仕様が明確でない場合(性能デバッグ等)は、システムを止めずに特定の状態へ絞り込める動的なログ収集・処理ツールとして働き、有用な述語を見つける探索を助ける。 - 性質の限界: 安全性のみ扱い、活性は、タイムアウトと安定性指標を述語に含めて誤警報を除く。デッドロック検出のように分割しにくい性質は、頻繁に状態が変わるロックを除く前段を足すなどの工夫が必要である。 - 部品レベルの述語は、複数システムが協調するデータセンターのアプリケーションでは不十分なことがある。予期しない相互作用の原因を特定できず、仕様が日々変わって最新でないためである。著者らは、この点と、オンライン検査とオフライン再生の統合を今後の課題とする。 ## 強み / 弱点・課題 - 強み - 述語を逐次コードで書ける低い導入コスト(最大の述語で 210 行、多くは約 100 行)。 - 無改造・実行中に述語を変更でき、常時稼働可能な低オーバーヘッド(最大 8%、多くは 1% 未満)。 - 大域スナップショットの構築と障害時の扱いが形式化され、検査側の障害耐性がある。 - 弱点・課題 - 有用な述語(仕様)が前提で、書けなければ効果が出ない。 - 大規模システムでの評価は検索エンジンの 250 サーバが最大で、Chord の実験は 85 ノードの検証環境である。 - 標本化と Tbuf の設定は、偽陰性や誤警報とのトレードオフを開発者が引き受ける。 - 状態露出は関数引数に限られる場合が多く、内部状態が引数に現れない系では C++ コードの埋め込みが要る。 - 挿入前の履歴に依存する違反は見逃す。安全性以外の性質、分割不能な述語(デッドロック等)には別の工夫を要する。 ## 関連 - ソース: [[@1978__CACM__Time Clocks and the Ordering of Events in a Distributed System]] - 概念: [[分散モニタリング]] / [[障害注入]] / [[分散コンセンサス]] / [[分散トレーシング]] - エンティティ: [[D3S]] / [[Xuezheng Liu]] / [[Zhenyu Guo]] / [[M. Frans Kaashoek]] / [[Microsoft Research Asia]] / [[MIT CSAIL]] / [[Tsinghua University]] / [[Fudan University]] / [[Shanghai Jiao Tong University]] ## 出典 - [[.raw/papers/2008__NSDI__D3S-Debugging-Deployed-Distributed-Systems.pdf]]