Contents
Apache Spark ストリーミング 入門 チュートリアル:ゼロから始めるハンズオンガイド
データ処理の現場では、リアルタイムでの分析ニーズが急増しています。そんな要望に応えるのが Apache Spark Streaming です。本記事では、初心者向けに 「ストリーミング処理の基本」をステップバイステップで解説し、実際のコードを通じて理解を深めます。特に、Java/Scala環境構築からワードカウントの実装まで、ハンズオン形式で進めます。
Apache Spark Streamingとは
リアルタイム処理に特化したApache SparkのモジュールであるSpark Streamingは、継続的なデータ流入をバッチ処理して分析する技術です。従来のバッチ処理と異なり、秒単位でデータを処理できることから、IoTやログ監視など幅広いシーンで活用されています。
リアルタイム処理の概要
Spark Streamingでは、ネットワークソケットやKafkaなどから継続的に流入するデータを「DStream(Discretized Stream)」という形式で扱います。このDStreamは時間単位(例: 1秒ごと)にバッチ化され、Spark CoreのRDD処理フローと統合されるため、高スループットかつ低レイテンシーを実現します。
Spark Coreとの関係性
Spark Streamingは、Spark Coreの上に構築された拡張機能です。バッチ処理で使うRDDやトランザクションセマンティクスなどの仕組みを活用し、リアルタイム処理を効率化します。この関係性により、既存のSpark知識がストリーミング開発にも流用可能です。
開発環境構築手順
実装するにはまず、Java/Scala環境とSparkを準備します。以下にLinux/macOS向けのセットアップ手順を示します。
Java/Scala環境準備
- Java:OpenJDK 8以上をインストール(
java -versionで確認) - Scala:2.12または2.13バージョンを導入(
scala -versionで確認) - IDE:IntelliJ IDEAやVS Codeなど、Scala対応のエディタを用意
blockquote: 環境構築に時間がかかる場合は、Dockerイメージを使用するのも一案です。
Sparkダウンロード・インストール
- Apache Spark公式サイトから最新バージョン(例: 3.5.x)をダウンロード
- 下記コマンドで解凍し、環境変数設定
bash
tar -xvf spark-3.5.0-bin-hadoop3.tgz
export SPARK_HOME=/path/to/spark
export PATH=$SPARK_HOME/bin:$PATH
SBTプロジェクト作成
SBT(Scala Build Tool)でプロジェクトを初期化し、依存関係を設定します。build.sbtに以下を記載:
|
1 2 3 4 5 6 7 8 |
name := "SparkStreamingTutorial" version := "0.1" scalaVersion := "2.13.10" libraryDependencies += "org.apache.spark" %% "spark-core" % "3.5.0" libraryDependencies += "org.apache.spark" %% "spark-streaming" % "3.5.0" libraryDependencies += "org.apache.spark" %% "spark-streaming-kafka-0-10" % "3.5.0" |
DStreamの基本概念
Spark Streamingの核となるのは、DStream(Discretized Stream)です。以下に仕組みと特徴を解説します。
Discretized Streamの仕組み
DStreamは、時間単位で区切られたデータの列として扱います。たとえば1秒ごとに流入するログデータは、「1秒分のバッチ」として処理され、RDDに変換されます。この「バッチ化」により、リアルタイム処理をバッチ処理の枠組みで実現します。
| トピック | 内容 |
|---|---|
| 時間粒度 | デフォルトは1秒(batchDurationで設定可能) |
| データ構造 | RDDのシーケンスとして管理 |
| 処理フロー | 受信 → バッチ化 → 変換 → 出力 |
トランザクションセマンティクス
DStreamは、失敗時の再試行やデータロスの防止に備えた「正確な実行」を担保します。この仕組みにより、リアルタイムアプリケーションでも信頼性が保たれます。
データソース接続方法
Spark Streamingでは、ネットワークまたはメッセージキューからデータを受け取る必要があります。代表的な2つの接続方法を解説します。
Socket接続の設定
テスト用にncコマンドで仮想のデータ送信環境を作成できます:
|
1 2 |
nc -lk 9999 |
このポートにデータを送信すると、Spark Streamingが受信します。コード例は後述のワードカウントで使用。
blockquote: Sparkアプリケーションを実行する前には、
ncコマンドでサーバーを起動し、別ターミナルでデータを送信してください。この手順が抜けていた場合、接続エラーが発生します。
Kafka統合の基本形
Kafkaとの接続には、spark-streaming-kafka-0-10ライブラリが必要です。以下のようにソースを作成:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 |
val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "localhost:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "spark-streaming-group", "auto.offset.reset" -> "latest" ) val locationStrategy = KafkaLocationStrategies.PreferConsistent val consumerStrategy = ConsumerStrategies.Subscribe[String, String](Seq("input-topic"), kafkaParams) val kafkaStream = KafkaUtils.createDirectStream[String, String]( locationStrategy, consumerStrategy, kafkaParams, topics = Seq("input-topic") ) |
ワードカウント処理の実装
ここでは、Socketから送信されたテキストをリアルタイムでワードカウントするコードを例に解説します。
コード構造の解説
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 |
import org.apache.spark._ import org.apache.spark.streaming._ object WordCount { def main(args: Array[String]) { val conf = new SparkConf().setAppName("WordCount") val ssc = new StreamingContext(conf, Seconds(1)) // Socketからデータを受信 val lines = ssc.socketTextStream("localhost", 9999) // ワードカウント処理(mapとreduceByKey) val words = lines.flatMap(_.split(" ")) val wordCounts = words.map(word => (word, 1)).reduceByKey(_ + _) // 結果をコンソール出力 wordCounts.print() ssc.start() ssc.awaitTermination() } } |
ステップバイステップ実行フロー
socketTextStreamでポート9999に接続- テキストをスペースごとに分割(
flatMap) - 各単語に対してカウントを1増やす(
map) - カウント値を集計(
reduceByKey) - 結果をコンソールに出力
blockquote: 実際のデータ送信は、別のターミナルで
nc -lk 9999を実行し、「hello world hello」などと入力してください。
実行結果の可視化と検証
処理が完了後、ローカル環境で確認する方法を解説します。
Spark Web UIの確認手順
- プログラムを実行中に、ブラウザで
http://localhost:4040にアクセス - Stagesタブ:処理のステップごとの進捗とリソース使用量が表示されます
- Executorsタブ:各ワーカーの状態を確認可能
カウンター値の検証方法
- コンソール出力に、カウント結果が秒単位で更新されるはずです。
text
Time: 2024-08-03 15:00:01
(hello,1)
(world,1)
...
まとめ
本記事では、Apache Spark Streamingの基本的な使い方をステップバイステップで解説しました。
- Spark Streamingとは:リアルタイム処理の仕組みとSpark Coreとの関係
- 環境構築方法:Java/Scala環境とSBTプロジェクトの作成手順
- DStreamの理解:データバッチングとトランザクションセマンティクス
- 実装例:Socket経由でのワードカウント処理コードと実行フロー
- 検証方法:Spark Web UIやコンソール出力による結果確認
記事内のサンプルコードをコピーしてローカル環境で実行し、リアルタイム処理の流れを体感してください。