このプロジェクトについて

このプロジェクトは、クラウドアカウントを一切必要とせず、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