这个项目能做什么
本项目提供了一套完全基于 Docker 的本地实时流式分析管道,无需任何云账户即可演示 Spark Structured Streaming 的核心概念。
## 项目功能
一个 Python 生产者会模拟点击流事件并发送至 Redpanda 主题(兼容 Kafka),而 PySpark Structured Streaming 作业则消费这些事件,对其应用带水印的 1 分钟滚动窗口聚合以处理迟到事件,并将结果输出到控制台或 Parquet 文件。
## 架构
- **Python Producer**:生成包含 URL 和时间戳数据的模拟点击流 JSON 事件
- **Redpanda**:在 Docker 中运行的与 Kafka 兼容的消息代理,无需 Zookeeper
- **PySpark Structured Streaming**:将 Kafka 读取为无界 DataFrame,使用显式 Schema 解析 JSON,应用带时间水印的事件时间窗口,并按 URL 和窗口输出聚合后的页面浏览量计数
## 演示的关键概念
- 使用 Spark Structured Streaming 将 Kafka 主题读取为无界 DataFrame
- 基于 Schema-on-Read 解析 JSON 事件载荷
- 带时间水印的事件时间窗口(1 分钟滚动窗口)以处理迟到的事件
- 按 URL 进行持续的微批次页面浏览量聚合
- 支持输出到控制台用于学习,并提供 Parquet 存储用于持久化
## 项目结构
```
streaming-pipeline-kafka-spark/
├── docker-compose.yml # Redpanda + Spark containers
├── producer/
│ └── click_event_producer.py # simulated clickstream producer
├── spark_app/
│ └── streaming_aggregation.py # Structured Streaming job
└── requirements.txt
```
## 快速开始
```bash
git clone https://github.com/Kornelius99/streaming-pipeline-kafka-spark.git
cd streaming-pipeline-kafka-spark
docker-compose up -d
# Terminal 1: start producing events
python producer/click_event_producer.py
# Terminal 2: run the Spark Streaming job
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
评论
0 评分人数达到10人后显示
登录后参与讨论。