这个项目能做什么

电商点击流分析平台是一个全栈数据工程项目,处理来自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