data-engineering-zoomcamp:用 Redpanda 与 PySpark Structured Streaming 构建实时出租车行程流式处理管线

📅 发布时间:2026/9/12 16:18:42
data-engineering-zoomcamp:用 Redpanda 与 PySpark Structured Streaming 构建实时出租车行程流式处理管线
data-engineering-zoomcamp用 Redpanda 与 PySpark Structured Streaming 构建实时出租车行程流式处理管线【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp本文以 Data Engineering Zoomcamp 2027 届 Streaming 模块中的 Redpanda PySpark 流式处理示例为核心完整讲解从基础设施准备、消息生产/消费到 Structured Streaming 实时聚合的端到端流程。读完本文你将掌握如何用 Docker 搭建 Redpanda 消息集群与 Spark 集群如何用 Kafka 客户端协议向 Redpanda 写入出租车行程 CSV 数据以及如何通过spark-submit提交 PySpark 流式计算任务对数据做实时分组与滑动窗口聚合。示例概览一条完整的实时数据处理链路该示例位于仓库 cohorts/2027/07-streaming/extras/python/streams-example/redpanda/整体数据流如下rides.csv ── producer.py ── Redpanda topic(rides_csv) ── consumer.py验证消费 └── streaming.pyPySpark 流式读取、解析、聚合 ├── Console sink实时查看 └── Kafka sink回写聚合结果 topic链路中的每一环都对应目录下的独立文件职责清晰文件职责docker-compose.yaml定义 Redpanda 集群含双节点配置示例与 Redpanda Console 管理界面settings.py集中管理数据路径、bootstrap 地址、topic 名称与 Spark 数据 Schemaproducer.py读取出租车行程 CSV以 key-value 形式写入 topicconsumer.py订阅 topic 验证消息已成功落盘支持--topic参数指定主题streaming.pyPySpark Structured Streaming 作业流式读取、Schema 解析、分组与窗口聚合、多种 Sink 输出spark-submit.sh提交脚本自动携带 Kafka 集成依赖包并提交 Spark 作业streaming-notebook.ipynb备选的 Notebook 交互式运行方式说明Redpanda 完全兼容 Kafka 协议因此示例中的生产者、消费者与 Spark 均使用 Kafka 客户端库无需任何 Redpanda 专有 SDK。第一步前置条件检查README 强调运行本示例前必须确保 Docker 网络与数据卷已经正确创建否则后续容器之间无法通信、共享工作目录也会失败。首先用以下命令确认docker volume ls # 应列出 hadoop-distributed-file-system docker network ls # 应列出 kafka-spark-network这两个资源在整个 Streaming 模块的多个示例Kafka、Redpanda、PySpark之间是共享的网络kafka-spark-network让 Redpanda 容器与 Spark 集群容器处于同一网络彼此通过容器名解析通信数据卷hadoop-distributed-file-system在容器与宿主机之间共享工作目录如数据文件与 checkpoint 目录。如果上面两条命令没有对应输出说明你还没有创建它们请执行第二步。第二步创建 Docker 网络与数据卷在尚未运行过其他示例的情况下需要手动创建网络与卷# 创建网络 docker network create kafka-spark-network # 创建数据卷 docker volume create --namehadoop-distributed-file-system网络与卷的具体挂载方式可以在 docker-compose.yaml 中看到shared-workspace卷的name被显式指定为hadoop-distributed-file-system网络则声明为external: true并命名为kafka-spark-network。external: true意味着该网络必须由外部预先创建这正是 README 要求先执行docker network create的原因volumes: shared-workspace: name: hadoop-distributed-file-system driver: local networks: default: name: kafka-spark-network external: true若你希望直接复用仓库中 Spark 集群的启动说明可参考 extras/python/docker/README.md其中的 Spark 集群编排JupyterLab、Spark Master、两个 Spark Worker同样依赖上述网络与卷具体见 docker/spark/docker-compose.yml。第三步启动 Redpanda 集群与 Console进入redpanda目录后启动服务docker compose up -d启动后你会得到三类容器1. Redpanda Brokerredpanda-1核心消息节点。关键启动参数如下command: - redpanda - start - --smp 1 # 仅使用 1 个 CPU 核心适合本地开发 - --reserve-memory 0M # 不预留额外内存降低本地资源占用 - --overprovisioned # 允许超卖适配开发环境 - --node-id 1 - --kafka-addr PLAINTEXT://0.0.0.0:29092,OUTSIDE://0.0.0.0:9092 - --advertise-kafka-addr PLAINTEXT://redpanda-1:29092,OUTSIDE://localhost:9092 - --pandaproxy-addr PLAINTEXT://0.0.0.0:28082,OUTSIDE://0.0.0.0:8082 - --advertise-pandaproxy-addr PLAINTEXT://redpanda-1:28082,OUTSIDE://localhost:8082 - --rpc-addr 0.0.0.0:33145 - --advertise-rpc-addr redpanda-1:33145这里需要重点理解的是Kafka 监听地址listener的 PLAINTEXT/OUTSIDE 双端口设计PLAINTEXT://redpanda-1:29092容器内部监听地址供同一 Docker 网络内的 Spark 等容器通过服务名redpanda-1访问OUTSIDE://localhost:9092宿主机访问地址供本机运行的producer.py、consumer.py通过localhost:9092连接与 settings.py 中的BOOTSTRAP_SERVERS localhost:9092一一对应Panda Proxy8082/28082提供 HTTP 访问入口rpc-addr33145用于节点间内部通信。2. 可选双节点集群redpanda-2默认随 compose 一起启动README 源码中将其注释为想要双节点集群取消注释即可其中通过--seeds redpanda-1:33145声明种子节点--node-id 2区分节点身份并使用独立的端口段9093/28083/33146避免冲突。本地学习单节点即可双节点主要用于观察集群扩展行为。3. Redpanda Consoleredpanda-console基于 Web 的可视化管理界面通过环境变量注入配置environment: CONFIG_FILEPATH: /tmp/config.yml CONSOLE_CONFIG_FILE: | kafka: brokers: [redpanda-1:29092] schemaRegistry: enabled: false redpanda: adminApi: enabled: true urls: [http://redpanda-1:9644] connect: enabled: false启动后可通过http://localhost:8080访问 Console在界面上查看 topic 列表、消息内容与消费进度是验证生产/消费结果最直观的手段。第四步认识配置中心 settings.pysettings.py 是整个示例的单一事实来源所有脚本共用同一套配置INPUT_DATA_PATH ../../resources/rides.csv # 输入数据相对 redpanda 目录的路径 BOOTSTRAP_SERVERS localhost:9092 # 宿主机视角的 Redpanda 地址 TOPIC_WINDOWED_VENDOR_ID_COUNT vendor_counts_windowed # 聚合结果回写 topic PRODUCE_TOPIC_RIDES_CSV CONSUME_TOPIC_RIDES_CSV rides_csv # 原始数据 topic RIDE_SCHEMA T.StructType([...]) # Spark 侧的行程数据 Schema其中RIDE_SCHEMA定义了 Spark 解析消息时的目标结构共 7 个字段vendor_idInteger、tpep_pickup_datetimeTimestamp、tpep_dropoff_datetimeTimestamp、passenger_countInteger、trip_distanceFloat、payment_typeInteger、total_amountFloat。该 Schema 将在第五步与生产者消息内容精确对应。第五步运行生产者 producer.py生产者使用 Python 版 Kafka 客户端kafka库即 kafka-python向 Redpanda 发送消息python producer.py从 producer.py 源码可以看到核心实现分为三部分1. 生产者初始化与序列化配置config { bootstrap_servers: [BOOTSTRAP_SERVERS], key_serializer: lambda x: x.encode(utf-8), value_serializer: lambda x: x.encode(utf-8) } producer RideCSVProducer(propsconfig)bootstrap_servers指向localhost:9092key 与 value 统一以 UTF-8 编码发送。2. CSV 读取与字段裁剪reader csv.reader(f) header next(reader) # 跳过表头 for row in reader: # vendor_id, passenger_count, trip_distance, payment_type, total_amount records.append(f{row[0]}, {row[1]}, {row[2]}, {row[3]}, {row[4]}, {row[9]}, {row[16]}) ride_keys.append(str(row[0])) i 1 if i 5: break代码从原始出租车行程 CSV 中抽取 7 列对应RIDE_SCHEMA的 7 个字段vendor_id、passenger_count、trip_distance、payment_type、total_amount以及 pickup/dropoff 时间戳列以vendor_id作为消息 key拼接成逗号分隔字符串作为 value。注意这里i 5即只读取前 5 条记录后退出——该示例刻意保持数据量小便于观察实际生产可移除该限制。3. 消息发送与 flushself.producer.send(topictopic, keykey, valuevalue) ... self.producer.flush()发送完成后调用flush()确保消息真正写入 broker而非仅停留在客户端缓冲区随后sleep(1)等待落盘完成。第六步运行消费者 consumer.py验证消息是否成功写入 topic 的快速方法# 使用默认 topicrides_csv运行 python consumer.py # 指定其他 topic python consumer.py --topic topic-nameconsumer.py 通过argparse暴露--topic参数默认值为配置中的CONSUME_TOPIC_RIDES_CSV即rides_csv。其消费配置同样值得逐项理解config { bootstrap_servers: [BOOTSTRAP_SERVERS], auto_offset_reset: earliest, # 从最早消息开始消费 enable_auto_commit: True, # 自动提交消费位点 key_deserializer: lambda key: int(key.decode(utf-8)), value_deserializer: lambda value: value.decode(utf-8), group_id: consumer.group.id.csv-example.1, }auto_offset_reset: earliest无已提交位点时从分区最老消息开始读保证不丢消息key_deserializerkey 反序列化为int对应生产端 key 为vendor_id数字字符串group_id消费组标识同组消费者共享消费位点。消费循环使用poll(1.0)以 1 秒为超时轮询源码注释说明轮询期间无法响应 SIGINT 中断故限制超时时间以便 CtrlC 退出逐条打印消息的 key 与 valueprint(fKey:{msg_val.key}-type({type(msg_val.key)}), fValue:{msg_val.value}-type({type(msg_val.value)}))预期输出形如Consuming from Kafka started Available topics to consume: {rides_csv} Key:1-type(class int), Value:1, 1, 1.1, 1, 17.3, 2021-01-01 00:31:03, 2021-01-01 00:33:03-type(class str)看到此类输出即代表生产→存储→消费链路已打通可以进入流式处理阶段。第七步运行 PySpark Structured Streaming 作业7.1 用 spark-submit.sh 提交README 明确说明spark-submit脚本会在运行streaming.py之前确保必要的 Kafka 集成 jar 包被安装。执行./spark-submit.sh streaming.pyspark-submit.sh 的实现要点if [ $# -lt 1 ]; then echo Usage: $0 pyspark-job.py [ executor-memory ] exit 1 fi PYTHON_JOB$1 EXEC_MEM${2:-1G} # 未指定第二个参数时默认 1G spark-submit --master spark://localhost:7077 --num-executors 2 \ --executor-memory $EXEC_MEM --executor-cores 1 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1,\ org.apache.spark:spark-avro_2.12:3.5.1,\ org.apache.spark:spark-streaming-kafka-0-10_2.12:3.5.1 \ $PYTHON_JOB脚本要求至少传入一个 PySpark 作业文件可选第二个参数指定 executor 内存字符串格式如512M或2G。它向运行在localhost:7077的 Spark Master 提交作业申请 2 个 executor、每 executor 1 核并通过--packages声明三个关键依赖spark-sql-kafka-0-10_2.12:3.5.1Structured Streaming 与 Kafka 协议集成的核心包spark-streaming-kafka-0-10_2.12:3.5.1旧版 Spark Streaming Kafka 集成用于兼容spark-avro_2.12:3.5.1Avro 序列化支持。运行前提Spark 集群Masterlocalhost:7077 至少 2 个 Worker已经启动。若使用仓库的 Docker 方案请先按 docker/spark/docker-compose.yml 启动 jupyterlab、spark-master、spark-worker-1、spark-worker-2 四个服务SPARK_WORKER_CORES1、SPARK_WORKER_MEMORY4g并确保 Redpanda 与 Spark 集群在同一kafka-spark-network中。7.2 streaming.py 的核心处理逻辑streaming.py 完整演示了 Structured Streaming 的读-解-算-写四步模型① 读取read_from_kafkadf_stream spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092,broker:29092) \ .option(subscribe, consume_topic) \ .option(startingOffsets, earliest) \ .option(checkpointLocation, checkpoint) \ .load()注意这里 bootstrap 地址同时给出了两个入口localhost:9092宿主机视角与broker:29092容器网络内视角这正对应 docker-compose 中 PLAINTEXT/OUTSIDE 双监听设计保证无论 Spark 运行在容器内还是宿主机都能连接。startingOffsets: earliest与消费者配置一致从最早消息开始处理。② 解析parse_ride_from_kafka_messageassert df.isStreaming is True, DataFrame doesnt receive streaming data df df.selectExpr(CAST(key AS STRING), CAST(value AS STRING)) col F.split(df[value], , ) for idx, field in enumerate(schema): df df.withColumn(field.name, col.getItem(idx).cast(field.dataType)) return df.select([field.name for field in schema])将 Kafka 消息的二进制 value 转为字符串后按生产端的分隔符,切分再依据RIDE_SCHEMA逐列cast成对应类型如tpep_pickup_datetime转 Timestamp。assert df.isStreaming在调试期保障调用方传入的确实是流式 DataFrame。③ 计算分组与滑动窗口聚合df_trip_count_by_vendor_id df.groupBy([vendor_id]).count() df_windowed_aggregation df.groupBy( F.window(timeColumndf.tpep_pickup_datetime, windowDuration10 minutes, slideDuration5 minutes), df.vendor_id ).count()示例并行演示两种聚合按vendor_id的全局计数以及基于tpep_pickup_datetime的10 分钟窗口、5 分钟滑动的行程计数即每 5 分钟输出一次最近 10 分钟的聚合结果。④ 输出三种 Sinksink_console(df_rides, output_modeappend) # 原始解析结果实时打印到控制台 sink_console(df_trip_count_by_vendor_id) # 分组计数打印到控制台默认 complete 模式 sink_kafka(dfdf_trip_count_messages, topicTOPIC_WINDOWED_VENDOR_ID_COUNT) # 窗口聚合回写 Kafkastreaming.py中提供了sink_console.format(console)可配processingTime触发间隔、sink_memory.format(memory)配合spark.sql查询、sink_kafka写回 topic需同时设置kafka.bootstrap.servers、topic、checkpointLocation三种输出方式。回写前通过prepare_df_to_kafka_sink将聚合结果拼成key/value两列key 取vendor_idvalue 取计数df df.withColumn(value, F.concat_ws(, , *value_columns)) if key_column: df df.withColumnRenamed(key_column, key) df df.withColumn(key, df.key.cast(string)) return df.select([key, value])最终spark.streams.awaitAnyTermination()让主进程持续运行等待所有流式查询终止。7.3 备选Notebook 交互式运行如果你更习惯 Notebook 环境仓库还提供了 streaming-notebook.ipynb。其首个代码单元通过环境变量注入依赖包os.environ[PYSPARK_SUBMIT_ARGS] --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.1,org.apache.spark:spark-avro_2.12:3.3.1 pyspark-shell从 Notebook 的已有运行记录execution_count输出可见Ivy 会解析并缓存spark-sql-kafka-0-10_2.12及其传递依赖spark-token-provider-kafka-0-10、kafka-clients、lz4-java、snappy-java等后续单元以与streaming.py相同的逻辑执行流式读取与聚合适合教学演示与逐步调试。注意 Notebook 中依赖版本为3.3.1而spark-submit.sh使用3.5.1二者需与本地 Spark 版本匹配。第八步验证与排错要点现象排查方向producer.py/consumer.py连接失败确认 Redpanda 容器已docker compose up -d确认docker network ls存在kafka-spark-networkdocker compose up报网络不存在未先执行docker network create kafka-spark-network按 README 顺序先创建网络与卷spark-submit找不到 Master确认 Spark Master 运行于localhost:7077Worker 数 ≥ 2消费者无输出但 Console 有消息检查--topic名称是否与PRODUCE_TOPIC_RIDES_CSVrides_csv一致流式作业未收到数据先运行python producer.py保证 topic 中有消息确认 Spark 端 bootstrap 地址容器内用broker:29092宿主机用localhost:9092频繁下载 jar首次spark-submit会经 Ivy 下载依赖属正常现象之后会命中本地缓存小结本示例用最小化的一组脚本串联起了数据生产 → 消息队列 → 流式消费 → 实时聚合 → 多路输出的完整实时数据管线基础设施层docker network/docker volume预先创建共享资源docker-compose.yaml 以 Kafka 协议兼容的 Redpanda 作为 broker并附 Console 可视化接入层producer.py 与 consumer.py 演示了标准 Kafka 客户端读写模式计算层streaming.py 展示了 Structured Streaming 的读-解-算-写全流程包括分组计数与 10 分钟/5 分钟滑动窗口聚合以及 Console 与 Kafka 两种 Sink 的搭配使用。该示例是理解消息队列 流式计算组合的经典起点——把rides.csv换成任意业务事件流把 Console Sink 换成文件、数据库或对象存储 Sink即可演变为生产级实时数据管线的雏形。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考