这个项目能做什么

本项目提供了一套完全基于 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