这个项目能做什么
电商点击流分析平台是一个全栈数据工程项目,处理来自Kaggle数据集2019年10月至11月期间约2.85亿条电商点击流事件。
## 概述
平台将原始CSV事件数据(约14 GB)摄入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 处理批处理ETL,包含五个PySpark作业:每日KPI计算、漏斗分析、购物车流失检测、商品关联度评分以及MapReduce的数据准备
- Apache Hadoop 3.4.1 在YARN上运行的Streaming执行两个MapReduce作业(品类事件计数和品牌收入聚合)
- Spark Structured Streaming 从Kafka消费数据,每5秒聚合活跃会话并写入PostgreSQL
**编排与API:**
- Apache Airflow 2.9.3 以DAG形式编排完整批处理管道
- FastAPI(配合asyncpg)为仪表板提供REST端点
- Terraform模板覆盖AWS和GCP部署
**前端:**
- React 18.3 结合TypeScript、Tailwind CSS和Recharts渲染四个页面:
- 概览:KPI卡片和每日趋势图表
- 漏斗分析:浏览到加购到购买的流程及各品类流失分析
- 商品分析:热门商品、品牌及品类
- 实时监控:通过Server-Sent Events实时追踪会话
## 数据管道
1. **批处理分析**:从MinIO读取原始CSV;计算每日指标(事件数、用户数、收入、转化率)、热门商品/品牌/品类并写入PostgreSQL。另一个作业计算各品类漏斗转化率、识别购物车流失,并计算商品关联提升得分(采样10%以保证性能)。
2. **MapReduce**:将准备好的Parquet文件转换为CSV并上传至HDFS,Hadoop Streaming作业按品类/事件类型计数事件并按品牌汇总收入。
3. **实时流处理**:Kafka生产者以可配置速率回放CSV(默认200条事件/秒)。Spark Streaming消费者每5秒微批次聚合活跃会话,并将其持久化至live_sessions表。
## 部署
全部13个服务通过Docker Compose运行。前置条件包括Docker Desktop >= 4.30(至少16 GB内存和30 GB磁盘空间)。Windows用户需WSL2且内存至少12 GB。Makefile提供了启动服务、上传数据、运行管道和执行测试的脚本。
启动后可访问的URL:仪表板(端口3000)、FastAPI文档(端口8000)、Airflow(端口8080)、MinIO控制台(端口9001)、Spark Master UI(端口8081)、HDFS NameNode(端口9870)及YARN(端口8088)。
## 仓库结构
关键目录包括`spark/jobs/`存放五个PySpark批处理脚本、`spark/streaming/`存放Kafka消费者、`kafka/producer/`存放回放生产者、`api/`存放带路由端点的FastAPI服务、`frontend/`存放React应用、`hadoop/mapreduce/`存放Streaming作业脚本、`airflow/dags/`存放编排逻辑、`terraform/`存放基础设施即代码模板,以及`docker/`存放Hadoop和Kafka服务定义。
许可证:MIT
评论
0 评分人数达到10人后显示
登录后参与讨论。