Skip to content

DNET_ApacheSpark

nishi_74322014 edited this page Sep 11, 2026 · 1 revision

Apache Spark

概要

Hadoop MapReduce(Hadoopの該当節を参照)と同様に

  • 複数の計算機を用いてデータ処理を行う並列データ処理系

  • JVM上で動作するOSSの並列分散処理系フレームワーク

    • 暗黙のデータ並列性と耐故障性を備えたクラスタ全体をプログラミングできる。
    • Resilient Distributed Dataset (RDD)と呼ばれる
      データ構造を処理するAPIを持つ。

背景

Apache Hadoop

Hadoop

  • データ処理してHDDに都度書き出す方式
  • ディスクIOを並列化してスループット高める。

Apache Spark(背景)

  • 大規模データの分散処理をオンメモリで実現する。

  • データ処理してHDDに都度書き出す方式よりも高速。

  • Hadoop MapReduce(Hadoopの該当節を参照)が
    適合しない以下のケースをサポートする。

    • 複雑なデータ処理を行うために,複数のジョブを連ねて実行する場合
    • 同じデータを複数のジョブから利用する場合
    • Hadoop Yarn(Hadoop)クラスタ上で動かすことも出来る。

特徴

分散処理

  • 以下の順序で、タスクにブレークダウンされる。

    • ジョブ

      • 一連のデータフロー(処理の全体)
      • 若しくは、SQLが生成した実行プラン
    • ステージ
      一纏まりの処理。

      • パーティション
        ・各ステージが処理する分割されたデータ
        ・パーティション数は、以下のいずれであるかによって決まる。
         ・データソース(例えば、HDFS(Hadoopの該当節を参照)なら128MB)
         ・中間データ(APIのパラメタや設定で決まる)
         ・Spark SQL(既定値200パーティションから自動シュリンク)

      • シャッフル
        ・ネットワーク越しのデータ転送を伴うデータの再分散。
        ・当該ステージのパーティションの出力を、
         次のステージのパーティションの入力にマップする。

      • タスク
        パーティションのデータをステージで処理する。

  • その他、関連する用語。

    • スロット
      タスクを割り当てるスロット(≒ CPUということらしい)。

    • 変換

      • ナロー変換
        クラスタネットワーク上でのデータシャッフルやデータ移動が不要な変換
      • ワイド変換
        クラスタネットワーク上でのデータシャッフルやデータ移動が必要な変換
    • パイプライン処理
      できる限り多くの操作をデータの 1 つのパーティションで実行すること。

      • データの 1 つのパーティションが RAM に読み取られると、
        1 つの タスク にできる限り多くのナロー操作が結合される。
      • ワイド操作では、シャッフルを強制するため、
        ステージを完了して、パイプラインを終了する。
  • 処理(ジョブ、ステージ、タスク)とデータ(RDD、パーティション)の時系列の関係
    (図:Apache Sparkのデータ処理の流れをなんとなく理解する - Qiita より引用)

SQLライク

Spark SQLによる、SQLライクな分散処理の隠蔽

トレードオフ

利点

  • 汎用的な並列データ処理系として利用できる。

    • RDDに対する数十種類のオペレータを利用可能。

      • 多様な並列データ処理をシンプルに記述できる。
      • オペレータを組み合わせれば、ジョブを組み合わせる必要がない。
      • 複数のオペレータは1つのタスクとしてRDDのパーティションごとにコピーされる。
      • パーティション並列性を活用でき、中間データI/Oを削減できる。
    • 標準で用途向けのライブラリが付属している。

  • 以下のようなビッグデータ シナリオに適合する。

欠点

苦手な処理。

  • クラスタ全体のメモリに乗り切らない 巨大なデータ処理(TB級以上)
  • 大きなデータセットを少しずつ更新する処理
  • 秒以下の特に短いレスポンスが必要な処理

シナリオ

適合するビッグデータ シナリオ

抽出、変換、読み込み (ETL)

  • Filtering
  • Sorting
  • Aggregating
  • Joining
  • Cleaning
  • Deduplicating
  • Validating

バッチ処理

Spark Streaming(シナリオ)

Spark Streaming

Spark SQL(シナリオ)

Spark SQL

用途向けのライブラリ(シナリオ)

その他、用途向けのライブラリ

詳細

アーキテクチャコンポーネントの関係が謎い
(詳細が見えて来たら書き足す予定)。

アーキテクチャ

ドライバ

  • プログラム
    コンソール アプリのようなプログラム

  • Spark セッション
    プログラムを受け取り、それを小さなタスクに分割する。
    小さくなったタスクはエグゼキュータで処理される。

エグゼキュータ

クラスタ

  • クラスタ
    エグゼキュータのホスティング

  • クラスター マネージャ
    次の目的でドライバエグゼキュータの両方と通信する。

    • リソースの割り当てを管理する
    • プログラム分割を管理する
    • プログラム実行を管理する

コンポーネント

Java、Scala、Python、R、C#、SQL
↓ ↓ ↓
Spark Streaming GraphX MLlib MLlib Structured Streaming
Spark SQL
Spark Core

Resilient Distributed Dataset (RDD)

Hadoop MapReduce(Hadoopの該当節を参照)が苦手としていた、
スループットとレイテンシの両立が必要な領域にアプローチ

  • 分散共有メモリを提供する分散プログラムのワーキングセット。

  • 永続化先として、主に、計算機のメモリ(キャッシュ)と二次記憶を利用できる。

  • メモリと二次記憶を組み合わせることも可能

    • メモリに保持しきれないパーティションを一時的に二次記憶に退避
    • 当該パーティションを利用する際に再び二次記憶から読み出す。

Spark Core

  • プロジェクト全体の基盤

  • RDDを抽象化した各種の実装

  • API(Java、Python、Scala、R)を介して公開

    • 分散タスクディスパッチ
    • スケジューリング
    • および基本I/O機能

Spark SQL

Spark SQL

  • DataFrameというオンメモリのテーブルに対し、
    LINQ的に処理を行うことが出来る。

  • 裏側では、
    クエリ・エンジン経由で実行プラン生成し

Spark Streaming

Spark Streaming

その他、用途向けのライブラリ

  • グラフ処理(GraphX)

  • 機械学習(Spark MLlib

  • 昨今はML Pipelinesの開発が活発
    scikit-learnのような
    機械学習全体のパイプラインをサポートするAPIが提供される。

耐障害性

≒ 再実行するタスクの数を最少にする機構。

ジョブ

物理ロギング(分散処理の該当節を参照)に依る。

  • 中間データをシャッフルする際に、
  • 中間データを二次記憶に書き出す。

RDD

以下の2つの方法に依る。

  • 再計算と呼ばれる論理ロギング(分散処理の該当節を参照)の一種
  • RDDの永続化と同時にレプリケーション(分散処理の該当節を参照)

CLI

spark-submit

Sparkアプリケーションを実行するコマンド

./bin/spark-submit \
  --class <main-class> \
  --master <master-url> \
  --deploy-mode <deploy-mode> \
  --conf <key>=<value> \
  ... # other options
  <application-jar> \
  [application-arguments]

, etc.

言語バインディング

基本はScalaで実装する。

PySpark

PySpark

.NET for Apache Spark

.NET for Apache Spark

チュートリアル

Apache Sparkチュートリアル

参考

分散処理 > 目的別

分散処理の該当節を参照。

分散(バッチ)系

分散処理の該当節を参照。

ストリーム系

分散処理の該当節を参照。

NTTデータ

先端技術株式会社

技術開発本部

https://www2.slideshare.net/nttdata-tech/presentations

システム技術本部

https://www.slideshare.net/hadoopxnttdata/presentations

Qiita

YARN

gihyo.jp … 技術評論社

Hadoopはどのように動くのか
─並列・分散システム技術から読み解くHadoop処理系の設計と実装

移行メモ

  • 「パーティション数はが何か?によって決まる。」は 「パーティション数は、以下のいずれであるかによって決まる。」に、 「限り多くの操作をデータの 1 つのパーティションで実行すること。」は 「できる限り多くの操作を〜」に、 「中間データを二時記憶に書き出す。」は「二次記憶」に、 「基本はScalarで実装する。」は「基本はScalaで実装する。」に、 「RDDのパーテョションごとに」は「パーティションごとに」に、 「[RDD]抽象化した各種の実装」は「RDDを抽象化した各種の実装」に正した。

  • 「分散処理」にあった処理とデータの関係の図は、元の PukiWiki で Qiita 上の外部画像を #ref で直接参照していたため、 画像の埋め込みではなく引用元記事へのリンクとした。

  • 「コンポーネント」の表は、元の PukiWiki で横結合(|>|)・縦結合(|~|)を 用いた積み上げ図だったため、結合部分を 〃 に置き換えた (MLlib が 2 列にあるのは RDD ベースと DataFrame ベースの 2 つを指すため原文どおり)。

  • マイクロソフト系技術情報 Wiki(techinfoofmicrosofttech.osscons.jp)への リンクは、移行済みの MS_DotNetForApacheSpark に張り替えた。


Tags: 移行, Apache Spark, 分散処理, RDD, Spark SQL, ビッグデータ, Hadoop

NetDevInfraWiki

マイクロソフト系技術情報 Wiki
Open 棟梁 Wiki

(未着手)

開発基盤部会 Wiki

移行管理: DONETODO

Clone this wiki locally