MapReduceは、クラスタ上で並列分散アルゴリズムを使用してビッグデータセットを処理および生成するためのプログラミングモデルおよび関連する実装です。[ 1 ] [ 2 ] [ 3 ]
MapReduceプログラムは、フィルタリングとソート(例えば、学生を名字順にキューにソートし、名前ごとに1つのキューを作成するなど)を実行するマッププロシージャと、集計操作(例えば、各キュー内の学生数をカウントし、名前の出現頻度を算出するなど)を実行するリデュースメソッドで構成されます。「MapReduceシステム」(「インフラストラクチャ」または「フレームワーク」とも呼ばれる)は、分散サーバーを統合し、さまざまなタスクを並列実行し、システムのさまざまな部分間のすべての通信とデータ転送を管理し、冗長性と耐障害性を提供することで、処理を統括します。
このモデルは、データ分析のための分割適用結合戦略の特殊化です。[ 4 ] これは、関数型プログラミングで一般的に使用されるmap 関数とreduce関数に触発されていますが、[ 5 ] MapReduce フレームワークにおけるそれらの目的は、元の形式とは異なります。[ 6 ] MapReduce フレームワークの主な貢献は、実際の map 関数と reduce 関数 (たとえば、1995 年のMessage Passing Interface標準の[ 7 ] reduce [ 8 ]およびscatter [ 9 ]操作に似ています) ではなく、並列化によってさまざまなアプリケーションで実現されるスケーラビリティとフォールトトレランスです。そのため、 MapReduce のシングルスレッド実装は通常、従来の (MapReduce ではない) 実装よりも高速ではありません。利点は通常、マルチプロセッサ ハードウェア上のマルチスレッド実装でのみ見られます。 [ 10 ]このモデルの使用は、MapReduceフレームワークの最適化された分散シャッフル操作(ネットワーク通信コストを削減する)とフォールトトレランス機能が適用される場合にのみ有益です。通信コストの最適化は、優れたMapReduceアルゴリズムにとって不可欠です。[ 11 ]
MapReduceライブラリは、さまざまな最適化レベルで、多くのプログラミング言語で記述されています。分散シャッフルをサポートする人気のオープンソース実装は、 Apache Hadoopの一部です。MapReduceという名前は元々Google独自の技術を指していましたが、その後、一般的な商標になりました。2014年までに、GoogleはもはやMapReduceを主要なビッグデータ処理モデルとして使用しておらず[ 12 ] 、 Apache Mahoutの開発は、完全なマップとリデュース機能を組み込んだ、より高性能でディスク指向性の低いメカニズムに移行しました[ 13 ] 。
MapReduceは、多数のコンピュータ(ノード)を使用して大規模データセットにわたる並列処理可能な問題を処理するためのフレームワークです。これらのコンピュータは、すべてのノードが同じローカルネットワーク上にあり、同様のハードウェアを使用している場合はクラスタ、地理的および管理的に分散したシステム間で共有され、より多様なハードウェアを使用している場合はグリッドと呼ばれます。処理は、ファイルシステム(非構造化データ)またはデータベース(構造化データ)に保存されたデータに対して実行できます。MapReduceはデータの局所性を活用し、データが保存されている場所の近くで処理することで、通信オーバーヘッドを最小限に抑えることができます。
MapReduceフレームワーク(またはシステム)は通常、次の3つの操作(またはステップ)で構成されます。
mapローカルデータに関数を適用し、出力を一時ストレージに書き込みます。マスターノードは、冗長な入力データのコピーが1つだけ処理されるようにします。map、1つのキーに属するすべてのデータが同じワーカーノード上に配置されるようにします。MapReduce は、マップ操作とリダクション操作の分散処理を可能にします。マップは、各マッピング操作が他の操作と独立している限り、並列で実行できます。実際には、これは独立したデータ ソースの数、および/または各ソースの近くにある CPU の数によって制限されます。同様に、マップ操作のすべての出力が同じキーを共有する場合、またはリダクション関数が連想的である場合、一連の「リデューサー」がリダクション フェーズを実行できます。このプロセスは、より逐次的なアルゴリズムと比較すると非効率に見えることが多いですが (リダクション プロセスの複数のインスタンスを実行する必要があるため)、MapReduce は、単一の「汎用」サーバーが処理できるよりもはるかに大きなデータセットに適用できます。大規模なサーバー ファームでは、MapReduce を使用して、わずか数時間で1ペタバイトのデータをソートできます。 [ 14 ]並列処理により、操作中にサーバーまたはストレージが部分的に故障した場合の復旧の可能性も提供されます。1 つのマッパーまたはリデューサーが故障した場合、入力データがまだ利用可能であると仮定して、作業を再スケジュールできます。
MapReduceを別の視点から捉えると、5段階の並列分散計算として理解できる。
これら5つのステップは論理的には順番に実行されるものと考えることができます。つまり、各ステップは前のステップが完了してから開始されます。ただし、最終結果に影響がない限り、実際にはこれらのステップを交互に実行することも可能です。
多くの場合、入力データは既に複数の異なるサーバーに分散(シャーディング)されている可能性があります。その場合、ステップ1では、ローカルに存在する入力データを処理するマップサーバーを割り当てることで、処理を大幅に簡略化できる場合があります。同様に、ステップ3では、処理する必要のあるマップ生成データにできるだけ近いリデュースプロセッサを割り当てることで、処理速度を向上させることができます。
MapReduceのMap 関数とReduce関数はどちらも、(キー、値)ペアで構造化されたデータに関して定義されています。Map関数は、あるデータドメインの型を持つ1組のデータを受け取り、別のドメインのペアのリストを返します。
Map(k1,v1)→list(k2,v2)
Map関数k1は、入力データセット内のすべてのペア(キーは)に並列に適用されます。これk2により、呼び出しごとにペア(キーは)のリストが生成されます。その後、MapReduceフレームワークは、すべてのリストから同じキー(k2)を持つすべてのペアを収集し、それらをグループ化して、キーごとに1つのグループを作成します。
次に、Reduce関数が各グループに並列に適用され、その結果、同じドメインの値の集合が生成されます。
Reduce(k2, list (v2))→ list((k3, v3))[ 15 ]
Reduce呼び出しは通常、1つのキーと値のペア、または空の戻り値を返しますが、1回の呼び出しで複数のキーと値のペアを返すことも可能です。すべての呼び出しの戻り値は、目的の結果リストとして収集されます。
このように、MapReduceフレームワークは(キー、値)ペアのリストを別の(キー、値)ペアのリストに変換します。[ 16 ]この動作は、任意の値のリストを受け取り、マップによって返されるすべての値を組み合わせた単一の値を返す、典型的な関数型プログラミングのマップとリデュースの組み合わせとは異なります。
MapReduceを実装するには、マップとリデュースの抽象化の実装が必要ですが、それだけでは十分ではありません。MapReduceの分散実装では、マップフェーズとリデュースフェーズを実行するプロセスを接続する手段が必要です。これは分散ファイルシステムでも構いません。他にも、マッパーからリデューサーへの直接ストリーミングや、マッピングプロセッサがクエリを実行するリデューサーに結果を提供するなど、様々な方法があります。
標準的な MapReduce の例では、一連のドキュメント内で各単語の出現回数をカウントします。[ 17 ]
function map (String name, String document): // name: ドキュメント名// document: ドキュメントの内容for each word w in document: 放出(w、1) function reduce (String word, Iterator partialCounts): // word: 単語// partialCounts: 集計された部分カウントのリスト 合計 = 0 partialCountsの各pcについて: 合計 += pc (単語、合計)を放出する
ここでは、各ドキュメントが単語に分割され、各単語は結果キーとしてマップ関数によってカウントされます。フレームワークは、同じキーを持つすべてのペアをまとめて、同じreduce関数呼び出しに渡します。したがって、この関数は入力値をすべて合計するだけで、その単語の出現回数を求めることができます。
別の例として、11億人のデータベースに対して、年齢別に一人当たりの平均ソーシャルコンタクト数を計算したい場合を考えてみましょう。SQLでは、このようなクエリは次のように表現できます。
SELECT age , AVG ( contacts ) FROM social.person GROUP BY age ORDER BY ageMapReduceを使用する場合、K1キーの値は1から1100までの整数で、それぞれが100万件のレコードのバッチを表し、K2キーの値は人の年齢(年)で、この計算は次の関数を使用して実行できます。
関数Map は、1 ~ 1100 の 整数K1 を入力として受け取ります。K1 バッチ内の各 social.person レコードに対して、 Y を人物の年齢とし 、Nを人物の連絡先の数とし、 出力レコード(Y,(N,1)) を 1 つ生成し、繰り返します。関数終了関数Reduceは入力として年齢 (年) Y を 受け取り、各入力レコード (Y,(N,C))に対して以下の処理を実行します。 SにN*Cの合計を蓄積するC new にC の 合計 を蓄積します。Aを S/C newとし、1 つの出力レコード(Y,(A,C new )) を生成します。関数終了
Reduce関数では、Cは合計N人の連絡先を持つ人の数を表すので、Map関数ではC=1と記述するのが自然です。なぜなら、各出力ペアは1人の人物の連絡先を参照しているからです。
MapReduce システムは 1100 個の Map プロセッサを並べ、それぞれに対応する 100 万個の入力レコードを提供します。Map ステップでは、例えば 8 から 103 の範囲のY値を持つ 11 億(Y,(N,1))個のレコードが生成されます。MapReduce システムは、年齢ごとの平均値が必要であるため、キー/値ペアのシャッフル操作を実行して 96 個の Reduce プロセッサを並べ、それぞれに対応する数百万個の入力レコードを提供します。Reduce ステップでは、Yでソートされた 96 個の出力レコード(Y,A)のみに大幅に削減されたセットが生成され、最終結果ファイルに格納されます。
レコード内のカウント情報は、処理が複数回削減される場合に重要です。レコード数を加算しないと、計算された平均値が誤ってしまいます。例えば、次のようになります。
-- マップ出力 #1: 年齢、接触回数 10、9 10、9 10、9
-- マップ出力 #2: 年齢、接触回数 10、9 10、9
-- マップ出力 #3: 年齢、接触回数 10、10
ファイル1とファイル2を削減すると、10歳の人の平均連絡先が9件の新しいファイルが作成されます((9+9+9+9+9)/5)。
-- ステップ1を減らす:年齢、平均接触者数 10、9
ファイル#3でこれを減らすと、既に見たレコードの数が失われるため、10 歳の平均コンタクト数が 9.5 ((9+10)/2) となり、これは間違っています。正しい答えは 9.1 66 = 55 / 6 = (9×3+9×2+10×1)/(3+2+1) です。
ソフトウェアフレームワークのアーキテクチャは、オープン/クローズドの原則に従っており、コードは変更不可能な固定部分と拡張可能なホットスポットに効果的に分割されています。MapReduceフレームワークの固定部分は、大規模な分散ソートです。アプリケーションが定義するホットスポットは次のとおりです。
入力リーダーは入力データを適切なサイズの「分割」(実際には通常64MB ~128MB )に分割し、フレームワークは各分割を各Map関数に割り当てます。入力リーダーは安定ストレージ(通常は分散ファイルシステム)からデータを読み込み、キーと値のペアを生成します。
一般的な例としては、テキストファイルが格納されたディレクトリを読み込み、各行をレコードとして返すというものがあります。
Map関数は、一連のキーと値のペアを受け取り、それぞれを処理して、0個以上の出力キーと値のペアを生成します。マップの入力型と出力型は、互いに異なる場合があり(そして多くの場合異なります)。
アプリケーションが単語数をカウントする場合、マップ関数は行を単語ごとに分割し、各単語に対応するキーと値のペアを出力します。各出力ペアには、キーとして単語、値として行内のその単語の出現回数が含まれます。
各Map関数の出力は、シャーディングの目的で、アプリケーションのパーティション関数によって特定のリデューサに割り当てられます。パーティション関数にはキーとリデューサの数が渡され、目的のリデューサのインデックスが返されます。
一般的なデフォルト設定では、キーをハッシュ化し、そのハッシュ値をリデューサーの数で割った余りを使用します。負荷分散のために、シャードごとにデータがほぼ均一に分散されるようなパーティション関数を選択することが重要です。そうしないと、MapReduce 操作が、処理の遅いリデューサー (つまり、不均一にパーティション化されたデータのより大きなシェアを割り当てられたリデューサー) の処理完了を待つために停止してしまう可能性があります。
マップ処理とリデュース処理の間では、データはシャッフル(並列ソート/ノード間での交換)され、データを生成したマップノードから、リデュース処理を行うシャードへと移動されます。ネットワーク帯域幅、CPU速度、生成されるデータ量、マップ処理とリデュース処理にかかる時間によっては、シャッフル処理に計算時間よりも時間がかかる場合があります。
各Reduceへの入力データは、 Mapが実行されたマシンから取得され、アプリケーションの比較関数を使用してソートされます。
このフレームワークは、ソートされた順序で一意のキーごとにアプリケーションのReduce関数を一度呼び出します。Reduce関数は、そのキーに関連付けられた値を反復処理し、0個以上の出力を生成します。
単語数をカウントする例では、Reduce関数は入力値を受け取り、それらを合計して、単語と最終的な合計値を単一の出力として生成します。
出力ライターは、R educeの出力を安定ストレージに書き込みます。
モノイドの特性は、MapReduce操作の妥当性を保証する基礎となる。[ 18 ] [ 19 ]
Algebird パッケージ[ 20 ]では、Map/Reduce の Scala 実装は明示的にモノイド クラス型を必要とします。[ 21 ]
MapReduceの操作は、2種類のデータタイプを扱います。1つはマッピングされる入力データのタイプA 、もう1つは削減される出力データのタイプBです。
Map操作は、型Aの個々の値を受け取り、各a:Aに対して値b:Bを生成します。Reduce操作は、型Bの値に対して定義されたバイナリ操作を必要とします。これは、利用可能なすべてのb:Bを単一の値に折り畳むことで構成されます。
基本的な要件の観点から言えば、MapReduce 操作には、削減対象のデータを任意に再グループ化する機能が含まれていなければなりません。このような要件は、操作の 2 つの特性に相当します。
2つ目の特性は、複数のノードで並列処理を行った場合、処理すべきデータを持たないノードは結果に影響を与えないことを保証するものです。
これら2つの性質は、演算・と中立元eを持つ型Bの値上のモノイド(B、・、e )を持つことに相当する。
型Aの値には制約はありません。任意の関数A → B をMap操作に使用できます。これは、カタモルフィズムA* → ( B、 •、e )が存在することを意味します。ここで、A*はクリーネ スターを表し、 A上のリストの型としても知られています。
シャッフル操作自体はMapReduceの本質とは関係ありません。クラウド上で計算を分散するために必要な操作です。
上記から、すべての二項リデュース操作がMapReduceで動作するわけではないことがわかる。以下に反例を示す。
MapReduce プログラムは必ずしも高速であるとは限りません。このプログラミング モデルの主な利点は、プラットフォームの最適化されたシャッフル操作を活用し、プログラムのMap 部分とReduce部分のみを記述すればよいことです。ただし、実際には、MapReduce プログラムの作成者はシャッフル ステップを考慮する必要があります。特に、パーティション関数とMap関数によって書き込まれるデータの量は、パフォーマンスとスケーラビリティに大きな影響を与える可能性があります。Combiner関数などの追加モジュールは、ディスクに書き込まれるデータ量とネットワーク経由で送信されるデータ量を削減するのに役立ちます。MapReduce アプリケーションは、特定の状況下では準線形の高速化を実現できます。[ 22 ]
MapReduceアルゴリズムを設計する際には、計算コストと通信コストの間で適切なトレードオフ[ 11 ]を選択する必要があります。通信コストは計算コストよりも大きいことが多く[ 11 ] [ 22 ]、多くのMapReduce実装ではクラッシュリカバリのためにすべての通信を分散ストレージに書き込むように設計されています。
MapReduceのパフォーマンスをチューニングする際には、マッピング、シャッフル、ソート(キーによるグループ化)、およびリデューシングの複雑さを考慮する必要があります。マッパーによって生成されるデータ量は、マッピングとリデューシングの間で計算コストの大部分をシフトさせる重要なパラメータです。リデューシングには、非線形な複雑さを持つソート(キーのグループ化)が含まれます。したがって、パーティションサイズが小さいほどソート時間は短縮されますが、リデューサーの数が多いと非現実的になる可能性があるため、トレードオフがあります。分割ユニットサイズの影響はごくわずかです(特に不適切な選択、たとえば1MB未満を選択しない限り)。一部のマッパーがローカルディスクから負荷を読み取ることによる利益は、平均的にはわずかです。[ 23 ]
処理が短時間で完了し、データが単一のマシンまたは小規模クラスタのメインメモリに収まる場合、MapReduceフレームワークの使用は通常効果的ではありません。これらのフレームワークは、計算中にノード全体が失われた場合でも復旧できるように設計されているため、中間結果を分散ストレージに書き込みます。このクラッシュリカバリはコストがかかり、計算に多数のコンピュータが関与し、実行時間が長い場合にのみ有効です。数秒で完了するタスクであれば、エラーが発生した場合でも単に再開するだけで済みますが、クラスタのサイズが大きくなるにつれて、少なくとも1台のマシンが故障する可能性は急速に高まります。このような問題では、すべてのデータをメモリに保持し、ノード障害が発生した場合に計算を再開する実装、またはデータが十分に小さい場合は非分散ソリューションの方が、MapReduceシステムよりも高速になることがよくあります。
MapReduceは、データセットに対する多数の操作をネットワーク内の各ノードに分割することで信頼性を実現します。各ノードは、完了した作業とステータス更新を定期的に報告することが求められます。ノードがその間隔よりも長く沈黙した場合、マスターノード(Googleファイルシステムのマスターサーバーに類似)はそのノードをダウンと記録し、そのノードに割り当てられた作業を他のノードに送信します。個々の操作では、並列で競合するスレッドが実行されていないことを確認するために、ファイル出力に名前を付けるアトミック操作を使用します。ファイルの名前を変更する際には、タスク名に加えて別の名前にコピーすることも可能です(副作用を許容するため)。
リデュース操作もほぼ同じように動作します。並列処理に対する特性が劣るため、マスターノードは、処理対象のデータを保持しているノードと同じノード、または同じラック内でリデュース操作を実行するようにスケジュールします。この特性は、データセンターのバックボーンネットワーク全体の帯域幅を節約できるため、望ましいものです。
実装は必ずしも高い信頼性を備えているとは限りません。例えば、Hadoopの旧バージョンでは、 NameNodeが分散ファイルシステムの単一障害点でした。Hadoopの最新バージョンでは、NameNodeのアクティブ/パッシブフェイルオーバーにより高可用性が実現されています。
MapReduce は、分散パターンベース検索、分散ソート、Web リンクグラフ反転、特異値分解[ 24 ]、 Web アクセスログ統計、転置インデックス構築、ドキュメントクラスタリング、機械学習[ 25 ]、統計的機械翻訳など、幅広いアプリケーションで役立ちます。さらに、MapReduce モデルは、マルチコアおよびメニーコア システム[ 26 ] [ 27 ] [ 28 ] 、デスクトップ グリッド [ 29 ]、マルチクラスタ[ 30 ] 、ボランティア コンピューティング環境[ 31 ] 、動的クラウド環境[ 32 ] 、モバイル環境[ 33 ] 、高性能コンピューティング環境[ 34 ]など、さまざまなコンピューティング 環境に適合されています。
Googleでは、MapReduceを使用してGoogleのワールドワイドウェブのインデックスを完全に再構築しました。これは、インデックスを更新し、さまざまな分析を実行する古いアドホックなプログラムに取って代わりました。 [ 35 ]その後、Googleの開発は、バッチ処理の代わりにストリーミング操作と更新を提供するPercolator、FlumeJava [ 36 ]、MillWheelなどのテクノロジーに移行し、完全なインデックスを再構築することなく「ライブ」検索結果を統合できるようにしました。[ 37 ]
MapReduceの安定した入力と出力は通常、分散ファイルシステムに保存されます。一時的なデータは通常、ローカルディスクに保存され、リデューサーによってリモートから取得されます。
並列データベースと共有なしアーキテクチャを専門とするコンピュータ科学者のDavid DeWitt 氏とMichael Stonebraker氏は、MapReduce が適用できる問題の範囲について批判的でした。[ 38 ]彼らは、そのインターフェースが低レベルすぎると述べ、支持者が主張するようなパラダイムシフトを本当に表しているのか疑問を呈しました。[ 39 ]彼らは、20 年以上前から存在する先行技術の例としてTeradataを挙げ、MapReduce の支持者の新規性に関する主張に異議を唱えました。また、MapReduce プログラマーとCODASYLプログラマーを比較し、どちらも「低レベルのレコード操作を実行する低レベル言語で記述している」と指摘しました。 [ 39 ] MapReduce は入力ファイルを使用し、スキーマをサポートしていないため、B ツリーやハッシュパーティショニングなどの一般的なデータベースシステムの機能によって実現されるパフォーマンスの向上を妨げていますが、 Pig (または PigLatin)、Sawzall、Apache Hive、[ 40 ] HBase [ 41 ]およびBigtable [ 41 ] [ 42 ]などのプロジェクトは、これらの問題の一部に対処しています。
グレッグ・ヨーゲンセンはこれらの見解を否定する記事を書いた。[ 43 ]ヨーゲンセンは、MapReduceはデータベースとして設計されたり使用されることを意図されたりしたことは一度もないため、デウィットとストーンブレーカーの分析全体は根拠がないと主張している。
DeWittとStonebrakerはその後、2009年にHadoopのMapReduceとRDBMSのアプローチのパフォーマンスをいくつかの特定の問題で比較した詳細なベンチマーク研究を発表しました。[ 44 ]彼らは、リレーショナルデータベースは、特に複雑な処理やデータが企業全体で使用されている場合など、多くの種類のデータ使用において真の利点を提供するが、MapReduceは単純な処理タスクや一度限りの処理タスクにはユーザーが採用しやすい可能性があると結論付けました。
MapReduceプログラミングパラダイムは、Danny Hillisの1985年の論文[ 45 ]でも記述されており、 Connection Machineでの使用を想定して「xapping/reduction」[ 46 ]と呼ばれ、マップとリデュースの両方を高速化するためにそのマシンの特殊なハードウェアに依存していた。最終的にConnection Machineで使用された方言である1986年のStarLispは並列処理*mapと[ 47 ]reduce!!を備えており、これはさらに非並列処理と組み込み機能を備えた1984年のCommon Lispに基づいていた[ 48 ] 。Connection Machineのハイパーキューブアーキテクチャが実行に用いるツリー状のアプローチは、mapreducereduce時間[ 49 ]は、Googleの論文で先行研究として言及されているアプローチと実質的に同じである。[ 3 ]:11
2010年、GoogleはMapReduceに関する特許を取得しました。2004年に出願されたこの特許は、Hadoop、CouchDBなどのオープンソースソフトウェアによるMapReduceの使用を対象としている可能性があります。Ars Technicaの編集者は、MapReduceの概念を普及させたGoogleの役割を認めつつも、この特許が有効か新規性があるか疑問を呈しました。[ 50 ] [ 51 ] 2013年、Googleは「オープン特許非主張(OPN)誓約」の一環として、この特許を防御的にのみ使用することを誓約しました。[ 52 ] [ 53 ]この特許は2026年12月23日に失効する予定です。[ 54 ]
MapReduceタスクは、非巡回データフロープログラム、つまりステートレスマッパーとステートレスリデューサーの順に記述し、バッチジョブスケジューラで実行する必要があります。このパラダイムでは、データセットの繰り返しクエリが困難になり、グラフ処理[ 55 ]などの分野で制限が生じます。グラフ処理では、単一の作業セットを複数回再訪する反復アルゴリズムが一般的です。また、レイテンシの高いディスクベースのデータが存在する場合、アルゴリズムは各パスでデータへのシリアルアクセスを許容できるにもかかわらず、データを複数回通過する必要がある機械学習の分野でも同様です[ 56 ] 。
のMapReduce実装は、どうやらMapReduceとはほとんど関係がないようです。私が読んだ限りでは、これはシングルスレッドで、MapReduceはクラスタ上で高度に並列化して使用されることを想定しているからです。... MongoDB MapReduceは、単一サーバー上でシングルスレッドで動作します...
「私たちはもはやMapReduceをほとんど使っていません」[Googleの技術インフラ担当上級副社長、ウルス・ヘルツレ氏]
氏のプレゼンテーションによると、10月時点でGoogleはMapReduceを通じて1日あたり約3,000件の計算ジョブを実行しており、これは数千マシン日に相当する。これらのバッチルーチンは、とりわけ最新のWebページを分析し、Googleのインデックスを更新する。