data-engineering-zoomcamp:用 Redpanda 与 PySpark Structured Streaming 构建实时出租车行程流式处理管线
2026/9/12 16:18:38 网站建设 项目流程

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.py(PySpark 流式读取、解析、聚合) ├──> Console sink(实时查看) └──> Kafka sink(回写聚合结果 topic)

链路中的每一环都对应目录下的独立文件,职责清晰:

文件职责
docker-compose.yaml定义 Redpanda 集群(含双节点配置示例)与 Redpanda Console 管理界面
settings.py集中管理数据路径、bootstrap 地址、topic 名称与 Spark 数据 Schema
producer.py读取出租车行程 CSV,以 key-value 形式写入 topic
consumer.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 --name=hadoop-distributed-file-system

网络与卷的具体挂载方式可以在 docker-compose.yaml 中看到:shared-workspace卷的name被显式指定为hadoop-distributed-file-system,网络则声明为external: true并命名为kafka-spark-networkexternal: 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 Broker(redpanda-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.pyconsumer.py通过localhost:9092连接(与 settings.py 中的BOOTSTRAP_SERVERS = 'localhost:9092'一一对应);
  • Panda Proxy(8082/28082)提供 HTTP 访问入口,rpc-addr33145)用于节点间内部通信。

2. 可选双节点集群(redpanda-2,默认随 compose 一起启动):README 源码中将其注释为"想要双节点集群?取消注释即可",其中通过--seeds redpanda-1:33145声明种子节点,--node-id 2区分节点身份,并使用独立的端口段(9093/28083/33146)避免冲突。本地学习单节点即可,双节点主要用于观察集群扩展行为。

3. Redpanda Console(redpanda-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.py

settings.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_id(Integer)、tpep_pickup_datetime(Timestamp)、tpep_dropoff_datetime(Timestamp)、passenger_count(Integer)、trip_distance(Float)、payment_type(Integer)、total_amount(Float)。该 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(props=config)

bootstrap_servers指向localhost:9092;key 与 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_idpassenger_counttrip_distancepayment_typetotal_amount,以及 pickup/dropoff 时间戳列),以vendor_id作为消息 key,拼接成逗号分隔字符串作为 value。注意这里i == 5即只读取前 5 条记录后退出——该示例刻意保持数据量小,便于观察;实际生产可移除该限制。

3. 消息发送与 flush

self.producer.send(topic=topic, key=key, value=value) ... self.producer.flush()

发送完成后调用flush()确保消息真正写入 broker(而非仅停留在客户端缓冲区),随后sleep(1)等待落盘完成。

第六步:运行消费者 consumer.py

验证消息是否成功写入 topic 的快速方法:

# 使用默认 topic(rides_csv)运行 python consumer.py # 指定其他 topic python consumer.py --topic <topic-name>

consumer.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_deserializer:key 反序列化为int(对应生产端 key 为vendor_id数字字符串);
  • group_id:消费组标识,同组消费者共享消费位点。

消费循环使用poll(1.0)以 1 秒为超时轮询(源码注释说明:轮询期间无法响应 SIGINT 中断,故限制超时时间以便 Ctrl+C 退出),逐条打印消息的 key 与 value:

print(f'Key:{msg_val.key}-type({type(msg_val.key)}), ' f'Value:{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.py

spark-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 内存(字符串格式如512M2G)。它向运行在localhost:7077的 Spark Master 提交作业,申请 2 个 executor、每 executor 1 核,并通过--packages声明三个关键依赖:

  • spark-sql-kafka-0-10_2.12:3.5.1:Structured Streaming 与 Kafka 协议集成的核心包;
  • spark-streaming-kafka-0-10_2.12:3.5.1:旧版 Spark Streaming Kafka 集成(用于兼容);
  • spark-avro_2.12:3.5.1:Avro 序列化支持。

运行前提:Spark 集群(Masterlocalhost:7077+ 至少 2 个 Worker)已经启动。若使用仓库的 Docker 方案,请先按 docker/spark/docker-compose.yml 启动 jupyterlab、spark-master、spark-worker-1、spark-worker-2 四个服务(SPARK_WORKER_CORES=1SPARK_WORKER_MEMORY=4g),并确保 Redpanda 与 Spark 集群在同一kafka-spark-network中。

7.2 streaming.py 的核心处理逻辑

streaming.py 完整演示了 Structured Streaming 的"读-解-算-写"四步模型:

① 读取:read_from_kafka

df_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_message

assert df.isStreaming is True, "DataFrame doesn't 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(timeColumn=df.tpep_pickup_datetime, windowDuration="10 minutes", slideDuration="5 minutes"), df.vendor_id ).count()

示例并行演示两种聚合:按vendor_id的全局计数,以及基于tpep_pickup_datetime10 分钟窗口、5 分钟滑动的行程计数(即每 5 分钟输出一次最近 10 分钟的聚合结果)。

④ 输出:三种 Sink

sink_console(df_rides, output_mode='append') # 原始解析结果实时打印到控制台 sink_console(df_trip_count_by_vendor_id) # 分组计数打印到控制台(默认 complete 模式) sink_kafka(df=df_trip_count_messages, topic=TOPIC_WINDOWED_VENDOR_ID_COUNT) # 窗口聚合回写 Kafka

streaming.py中提供了sink_console.format("console"),可配processingTime触发间隔)、sink_memory.format("memory"),配合spark.sql查询)、sink_kafka(写回 topic,需同时设置kafka.bootstrap.serverstopiccheckpointLocation)三种输出方式。回写前通过prepare_df_to_kafka_sink将聚合结果拼成key/value两列,key 取vendor_id,value 取计数:

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-10kafka-clientslz4-javasnappy-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-network
docker compose up报网络不存在未先执行docker network create kafka-spark-network,按 README 顺序先创建网络与卷
spark-submit找不到 Master确认 Spark Master 运行于localhost:7077,Worker 数 ≥ 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询