ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Data Engineering Zoomcamp 实战:使用 Redpanda 运行 PySpark Structured Streaming 流式管道

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 仓库中 redpanda 流式示例完整演示如何用 Docker 搭建 RedpandaKafka 协议兼容的消息队列通过 Python Producer 发布出租车行程Ride数据再用 PySpark Structured Streaming 实时消费、聚合并回写结果主题。读完本文你将掌握从环境初始化、数据生产、流式读取、Schema 解析到窗口聚合与多路 Sink 的一整套可运行的流式处理方案。1. 示例整体架构该示例位于07-streaming/extras/python/streams-example/redpanda/目录是一个完整的数据入队 → 实时消费 → 流式计算 → 结果落盘闭环Redpanda 集群由 docker-compose.yaml 启动对外暴露 Kafka 协议端口9092Producerproducer.py读取本地 CSV 数据按行发布到rides_csv主题Consumerconsumer.py订阅主题验证消息的 key/value 是否按预期到达Streaming Jobstreaming.pyPySpark Structured Streaming 消费rides_csv解析字段、按 vendor_id 计数、做 10 分钟窗口5 分钟滑动聚合结果同时输出到控制台和vendor_counts_windowed主题Spark 提交脚本spark-submit.sh提交任务前自动拉取 Kafka/Avro 连接器依赖。2. 前置条件网络与共享卷Redpanda 容器与 Spark 作业之间通过共享卷交换数据示例默认约定名为hadoop-distributed-file-system的 Docker volume名为kafka-spark-network的 Docker network。开始前先确认二者是否已存在docker volume ls # 应列出 hadoop-distributed-file-system docker network ls # 应列出 kafka-spark-network若未创建按下面的命令补齐。3. 创建 Docker 网络与数据卷如果此前未运行过本仓库其他示例ls没有任何输出则执行# 创建网络 docker network create kafka-spark-network # 创建数据卷 docker volume create --namehadoop-distributed-file-system3.1 卷与网络的落地方式在 docker-compose.yaml 中可以看到它们如何被引用volumes: shared-workspace: name: hadoop-distributed-file-system driver: local networks: default: name: kafka-spark-network external: true卷以shared-workspace为服务内别名实际物理名称固定为hadoop-distributed-file-system这样 Spark 侧与 Redpanda 侧可以挂载同一份工作目录网络声明为external: true即复用第 3 节手工创建的网络保证 Redpanda 集群与后续 Spark 容器处于同一网段。4. 启动 Redpanda 集群示例提供了一节点默认与两节点两种拓扑docker compose up -dredpanda-1默认启用镜像docker.redpanda.com/redpandadata/redpanda:v23.2.26对外暴露 Kafka 端口9092OUTSIDE 监听同时暴露 PandaProxy HTTP 端口8082与管理端口9644redpanda-2默认注释按需取消注释镜像redpanda:v23.1.1通过--seeds redpanda-1:33145加入集群对外端口9093用于体验多节点集群redpanda-console镜像console:v2.2.2提供 Web UIhttp://localhost:8080其kafka.brokers指向容器内地址redpanda-1:29092。关键启动参数说明以 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:9092OUTSIDE://localhost:9092是宿主机侧 Producer/Consumer/Spark 使用的地址这也解释了为何 Python 脚本与spark-submit.sh中统一使用localhost:9092作为 bootstrap server。5. 运行 Producer 与 Consumer5.1 Producer把 CSV 行程数据发布到 Kafkapython producer.pyproducer.py 的核心逻辑读取数据read_records()打开 rides.csv路径../../resources/rides.csv由 settings.py 中的INPUT_DATA_PATH定义跳过表头后取出每行的指定列拼成字符串vendor_id, tpep_pickup_datetime, tpep_dropoff_datetime, passenger_count, trip_distance, payment_type, total_amount默认只取前 5 条用于演示序列化与发送配置key_serializer/value_serializer为 UTF-8 编码通过KafkaProducer.send(topic, key, value)异步发送随后flush()确保落盘目标主题PRODUCE_TOPIC_RIDES_CSV CONSUME_TOPIC_RIDES_CSV rides_csv。以 CSV 首行为例实际发送的 value 形如1, 2020-07-01 00:25:32, 2020-07-01 00:33:39, 1, 1.50, 2, 9.35.2 Consumer验证消息到达# 使用默认设置运行消费者 python consumer.py # 指定消费特定主题 python consumer.py --topic topic-nameconsumer.py 通过 argparse 暴露--topic参数默认值为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, }consume_from_kafka()使用subscribe()订阅主题并poll(1.0)轮询代码注释特别说明SIGINT 无法在 poll 期间被处理因此把超时限制为 1 秒保证 CtrlC 能及时退出并执行consumer.close()。运行后应看到每条记录的Key整型 vendor_id与Value行程字段串。6. 运行 Streaming 脚本PySpark Structured Streaming6.1 一键提交脚本./spark-submit.sh streaming.pyspark-submit.sh 在真正提交前自动完成依赖安装避免手动拼--packagesPYTHON_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提交到本地 Standalone 集群spark://localhost:70772 个 executor、每 executor 1 核依赖锁定为 Spark 3.5.1 对应的 Kafka 0-10 集成包与 Avro 包如需调整内存可传第二个参数如./spark-submit.sh streaming.py 2G。6.2 Streaming Job 的完整链路streaming.py 按读 → 解析 → 聚合 → 多路输出组织第一步读取 Kafka 数据流def read_from_kafka(consume_topic: str): 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() return df_stream注意 bootstrap 地址同时包含宿主机地址localhost:9092与容器内地址broker:29092可兼容不同网络位置。第二步解析 value 并映射为强类型 SchemaProducer 写入的是逗号分隔的字符串因此先CAST(value AS STRING)再按, 切分为数组依据 settings.py 中的RIDE_SCHEMA逐列cast成对应类型RIDE_SCHEMA T.StructType([ T.StructField(vendor_id, T.IntegerType()), T.StructField(tpep_pickup_datetime, T.TimestampType()), T.StructField(tpep_dropoff_datetime, T.TimestampType()), T.StructField(passenger_count, T.IntegerType()), T.StructField(trip_distance, T.FloatType()), T.StructField(payment_type, T.IntegerType()), T.StructField(total_amount, T.FloatType()), ])parse_ride_from_kafka_message()用assert df.isStreaming is True校验输入确为流式 DataFrame随后withColumn(field.name, col.getItem(idx).cast(field.dataType))按索引展开成 7 个顶层列——这正是字段顺序必须与 Producer 拼接顺序一致的原因。第三步三种聚合与多路 Sink# 1) 原样输出到控制台append 模式5 秒触发一次 sink_console(df_rides, output_modeappend) # 2) 按 vendor_id 分组计数complete 模式输出到控制台 df_trip_count_by_vendor_id op_groupby(df_rides, [vendor_id]) # 3) 10 分钟窗口、5 分钟滑动的计数 df_windowed op_windowed_groupby( df_rides, window_duration10 minutes, slide_duration5 minutes )窗口聚合使用F.window(timeColumndf.tpep_pickup_datetime, windowDuration..., slideDuration...)即基于tpep_pickup_datetime事件时间而非到达时间做水印计算。第四步结果回写 Kafka 主题df_trip_count_messages prepare_df_to_kafka_sink( dfdf_trip_count_by_pickup_date_vendor_id, value_columns[count], key_columnvendor_id ) kafka_sink_query sink_kafka( dfdf_trip_count_messages, topicTOPIC_WINDOWED_VENDOR_ID_COUNT # vendor_counts_windowed )prepare_df_to_kafka_sink()用concat_ws(, , *value_columns)把 count 拼成 value、把 vendor_id 重命名为 key 并 cast 为字符串sink_kafka()以complete模式写入vendor_counts_windowed主题从而形成消费 → 计算 → 再生产的流式闭环。最后spark.streams.awaitAnyTermination()挂起主线程等待所有 StreamingQuery 结束。7. 从 notebook 交互式体验除脚本外目录下还提供 streaming-notebook.ipynb适合在 Jupyter 中以 Cell 为单位逐步执行读流、解析、聚合与查询便于调试窗口参数与输出模式console / memory sink。sink_memory()将结果注册为内存临时表之后可用spark.sql()直接查询是验证聚合结果的轻量手段。8. 常见排查点与使用前提先启动 Redpanda 再跑脚本Producer/Consumer 依赖localhost:9092若未先docker compose up -d会报连接拒绝共享卷/网络必须存在docker compose up会因external: true网络不存在而失败务必先执行第 3 节命令端口冲突9092、8082、9644等端口若被占用需调整 compose 中的映射Spark 集群地址spark-submit.sh写死spark://localhost:7077本地未起 Spark Standalone 时请先启动集群或按需修改 master 地址版本前提本示例面向 Spark 3.5.1 与 Redpanda v23.2.26两节点示例为 v23.1.1依赖坐标spark-sql-kafka-0-10_2.12:3.5.1等与之一一对应。9. 延伸阅读本示例与仓库内其他流式模块可相互印证07-streaming 课程主目录覆盖 Redpanda、Python 生产/消费、Flink 等完整流式知识体系pyspark 同构示例同样的 Producer/Consumer/Streaming 三件套结构便于对比不同主题接入方式Kafka 主题消费示例展示 JSON 格式的轻量生产/消费写法官方 Structured Streaming 编程指南与 Kafka 集成指南可作原理性参考本仓库实践基于上述文件路径内的实现为准。【免费下载链接】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),仅供参考
返回列表