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

E-Commerce Clickstream Analytics Platformは、2019年10〜11月のKaggleデータセットから約2.85億件のECクリックストリームイベントを処理するフルスタックデータエンジニアリングプロジェクトです。 ## 概要 本プラットフォームは、生のCSVイベントデータ(約14GB)をS3互換オブジェクトストア(MinIO)に取り込み、Apache Sparkによるバッチ分析、KafkaおよびSpark Structured Streamingによるライブストリーミング、Hadoop MapReduceジョブの実行を行い、FastAPIバックエンドを通じてReactフロントエンドダッシュボードに結果を配信します。 ## アーキテクチャコンポーネント **取り込み&ストレージ:** - MinIO(S3互換オブジェクトストレージ)が生のCSVおよびParquet出力を保持 - PostgreSQL 16.3が集約KPIおよびクエリ対応結果を保存 - Redis 7.4.1がキャッシュレイヤーを提供 - Apache Kafka 3.9(KRaftモード)がストリーミングデータのメッセージバスとして機能 **処理:** - Apache Spark 3.5が5つのPySparkジョブによるバッチETLを処理:日次KPI計算、ファンネル分析、カート放棄検出、商品アフィニティスコアリング、MapReduce用データ準備 - Apache Hadoop 3.4.1(YARN Streaming実行)が2つのMapReduceジョブを実行(カテゴリイベントカウント、ブランド収益集計) - Spark Structured StreamingがKafkaから消費し、5秒ごとにアクティブセッションを集計してPostgreSQLにupsert **オーケストレーション&API:** - Apache Airflow 2.9.3がDAGとして全体のバッチパイプラインをオーケストレート - FastAPI(asyncpg付き)がダッシュボード用のRESTエンドポイントを公開 - TerraformテンプレートがAWSおよびGCPデプロイメントをカバー **フロントエンド:** - React 18.3(TypeScript、Tailwind CSS、Recharts)により4ページをレンダリング: - Overview:KPIカードおよび日次トレンドチャート - Funnel Analysis:ビューからカート購入へのフローおよびカテゴリ別放棄率 - Product Analytics:上位商品・ブランド・カテゴリ - Live Monitor:Server-Sent Eventsによるリアルタイムセッション追跡 ## データパイプライン 1. **バッチ分析**:生のCSVをMinIOから読み込み、日次メトリクス(イベント数、ユーザー数、収益、コンバージョン率)、上位商品/ブランド/カテゴリを計算してPostgreSQLに書き込みます。別のジョブでカテゴリ別ファンネル率を計算し、カート放棄を特定し、商品アフィニティ_lift_スコア(パフォーマンス用に10%サンプリング)を算出します。 2. **MapReduce**:用意したParquetファイルをCSVに変換してHDFSにアップロードし、Hadoop Streamingジョブでカテゴリ/イベントタイプ別のイベントカウントおよびブランド別収益合計を実行します。 3. **ライブストリーミング**:KafkaプロデューサーがCSVを設定可能レート(デフォルト200イベント/秒)でリプレイします。Spark Streamingコンシューマーが5秒間のマイクロバッチでアクティブセッションを集計し、`live_sessions`テーブルに永続化します。 ## デプロイメント 全13サービスのすべてがDocker Composeで実行されます。前提条件としてDocker Desktop >= 4.30(16GB RAM、30GBディスクスペース以上)が必要です。WindowsユーザーはWSL2と12GB以上のメモリを必要とします。Makefileには、サービス起動、データアップロード、パイプライン実行、テスト実行のターゲットが含まれています。 起動後のアクセスURL:ダッシュボード(ポート3000)、FastAPI docs(ポート8000)、Airflow(ポート8080)、MinIO Console(ポート9001)、Spark Master UI(ポート8081)、HDFS NameNode(ポート9870)、YARN(ポート8088)。 ## リポジトリ構造 主要ディレクトリ:`spark/jobs/`(5つのPySparkバッチスクリプト)、`spark/streaming/`(Kafkaコンシューマー)、`kafka/producer/`(リプレイプロデューサー)、`api/`(ルーティング済みエンドポイントを持つFastAPIサービス)、`frontend/`(Reactアプリケーション)、`hadoop/mapreduce/`(Streamingジョブスクリプト)、`airflow/dags/`(オーケストレーション)、`terraform/`(IaCテンプレート)、`docker/`(HadoopおよびKafkaサービス定義) ライセンス:MIT