このプロジェクトについて
このプロジェクトは、クラウドアカウントを一切必要とせず、Spark Structured Streamingの概念を示す完全にDocker化されたローカルなリアルタイムストリーミング分析パイプラインを提供します。
## 主な機能
PythonプロデューサーがRedpandaトピック(Kafka互換)にクリックストリームイベントを送信し、PySpark Structured Streamingジョブはそのイベントを消費して、遅延イベント用のウォーターマーキング付き1分間のトゥーリングウィンドウ集約を適用し、結果をコンソールまたはParquetファイルに出力します。
## アーキテクチャ
- **Pythonプロデューサー**: URLとタイムスタンプデータを含むシミュレートされたクリックストリームJSONイベントを生成
- **Redpanda**: Zookeeperを必要とせずDockerで実行されるKafka互換ブローカー
- **PySpark Structured Streaming**: KafkaからアンバウンドDataFrameとして読み取り、明示的なスキーマでJSONを解析し、イベントタイムウィンドウイングとウォーターマーキングを適用し、ウィンドウごとのURL別のページビューカウントを集約して出力
## 実装されている主要コンセプト
- Spark Structured StreamingによるKafkaトピックをアンバウンドDataFrameとして読み込み
- スキーマオンリードによるJSONイベントペイロードの解析
- 遅延イベント対応のウォーターマーキング付きイベントタイムウィンドウイング(1分間トゥーリングウィンドウ)
- URLごとのページビュー継続マイクロバッチ集約
- 学習用コンソール出力、および永続化用Parquetシンクオプション
## プロジェクト構造
```
streaming-pipeline-kafka-spark/
├── docker-compose.yml # Redpanda + Sparkコンテナ
├── producer/
│ └── click_event_producer.py # シミュレートされたクリックストリームプロデューサー
├── spark_app/
│ └── streaming_aggregation.py # Structured Streamingジョブ
└── requirements.txt
```
## 始め方
```bash
git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git
cd streaming-pipeline-kafka-spark
docker-compose up -d
# ターミナル1: イベント送信を開始
python producer/click_event_producer.py
# ターミナル2: Spark Streamingジョブを実行
docker-compose exec spark spark-submit \
--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 \
/opt/spark_app/streaming_aggregation.py
```
結果は1分単位のウィンドウごとのページビューカウントとしてコンソールに表示され、新しいイベントが届くたびに更新されます。
## なぜRedpandaか
RedpandaはKafkaプロトコルをサポートしているため既存のKafkaクライアントコードは変更なしで動作しますが、Zookeeper依存がなくシングルコンテナで提供されるため、docker-composeをシンプルに保ち学習目的に適しています。
## ライセンス
MIT
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.