ARTICLE DETAIL

资讯详情

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

Data Engineering Zoomcamp 实战:使用 PySpark Structured Streaming 消费 Kafka 流并完成窗口聚合

Data Engineering Zoomcamp 实战:使用 PySpark Structured Streaming 消费 Kafka 流并完成窗口聚合 Data Engineering Zoomcamp 实战使用 PySpark 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本篇技术指南基于 Data Engineering Zoomcamp 仓库中 pyspark 流处理示例完整讲解如何用 PySpark Structured Streaming 从 Kafka 读取出租车行程数据、按 Schema 解析二进制消息、执行分组与滑动窗口聚合并把结果写回 Kafka。读完本文你将掌握一套可复用的「Kafka Producer → PySpark Streaming 消费 → Console / Memory / Kafka Sink」端到端流水线的搭建与运行方法。前置条件先让 Kafka 与 Spark 服务跑起来PySpark 流处理示例依赖 Kafka 与 Spark 两个集群。仓库在 docker 目录 中提供了完整的 Docker 编排方案包含 Kafka、Schema Registry、Zookeeper、Control Center以及 JupyterLab、Spark Master、Spark Worker 等容器。启动前必须先创建共享 Volume与外部网络这是整个示例能否联通的先决条件。README 给出了两个校验命令docker volume ls # should list hadoop-distributed-file-system docker network ls # should list kafka-spark-networkVolume 与网络的具体创建方式在 docker README 中有详细说明# 创建网络 docker network create kafka-spark-network # 创建共享卷Spark 集群文件系统 docker volume create --namehadoop-distributed-file-system为什么这两个资源如此重要从 顶层 docker-compose.yml 可以看到Spark 集群的 docker-compose.yml 中volumes.shared-workspace的name被显式指定为hadoop-distributed-file-system让 JupyterLab、Spark Master 与所有 Worker 挂载同一份工作目录/opt/workspace提交作业时代码与数据才能互相可见networks.default被设置为name: kafka-spark-network且external: true声明该网络是预先手动创建的Spark 容器加入后即可与 Kafka 的broker容器互通域名。如果跳过创建步骤docker compose up会因找不到外部网络而直接失败。创建好资源后在kafka与spark两个子目录分别执行docker compose up -d即可拉起全部服务。值得一提的是Kafka 容器在 docker-compose.yml 中配置了双监听器LISTENER_BOB://broker:29092供容器内互通LISTENER_FRED://localhost:9092供宿主机访问。这也解释了后文代码中为何 bootstrap servers 总是写成localhost:9092,broker:29092的混合形式。全局配置settings.py 与流 Schemapyspark 示例目录 下的所有脚本共用同一个 settings.py 作为配置中心import pyspark.sql.types as T INPUT_DATA_PATH ../../resources/rides.csv BOOTSTRAP_SERVERS localhost:9092 TOPIC_WINDOWED_VENDOR_ID_COUNT vendor_counts_windowed PRODUCE_TOPIC_RIDES_CSV CONSUME_TOPIC_RIDES_CSV rides_csv 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()), ])关键配置项一览配置项值作用INPUT_DATA_PATH../../resources/rides.csv生产者读取的出租车行程 CSV 路径BOOTSTRAP_SERVERSlocalhost:9092Python 生产者/消费者访问 Kafka 的地址PRODUCE_TOPIC_RIDES_CSV/CONSUME_TOPIC_RIDES_CSVrides_csv行程原始数据所在的 Kafka 主题TOPIC_WINDOWED_VENDOR_ID_COUNTvendor_counts_windowed聚合结果回写 Kafka 的目标主题RIDE_SCHEMA7 个字段的 StructType定义流数据的字段名与类型RIDE_SCHEMA是流处理的核心契约它决定了 Spark 如何把 Kafka 中逗号分隔的字符串消息解析成结构化列。字段类型的选择也直接对应 Kafka 消息的实际格式vendor_id、passenger_count、payment_type为整数tpep_pickup_datetime、tpep_dropoff_datetime为时间戳trip_distance、total_amount为浮点数。运行生产者与消费者打通 Kafka 数据链路生产者 producer.pyproducer.py 使用kafka-python库把 CSV 文件中的行程记录发送到rides_csv主题。它定义了一个RideCSVProducer类read_records()读取 CSV跳过表头后仅抽取 7 列row[0]、row[1]、row[2]、row[3]、row[4]、row[9]、row[16]重新拼接为逗号分隔字符串并将vendor_id作为消息 key示例代码中为演示只读取前 5 条记录publish()通过producer.send(topic, key, value)逐条发送最后flush()确保消息落盘序列化器把 key 与 value 统一encode(utf-8)。启动生产者python3 producer.py运行后可看到类似输出Producing record for key: 1, value:1, 2020-07-01 00:25:32, 2020-07-01 00:33:39, 1, 1.5, 2, 9.3消费者 consumer.pyconsumer.py 实现了一个基于轮询的RideCSVConsumer通过--topic参数指定消费主题默认值为rides_csv对应CONSUME_TOPIC_RIDES_CSV配置了auto_offset_resetearliest从头开始消费、enable_auto_commitTrue、group id 为consumer.group.id.csv-example.1key 反序列化为整数int(key.decode(utf-8))value 保持字符串poll(1.0)每次轮询 1 秒收到消息即打印 key/value 及类型。运行消费者验证数据链路# 默认主题 python3 consumer.py # 指定主题 python3 consumer.py --topic topic-name注意生产者与消费者脚本运行于宿主机而非容器因此它们使用localhost:9092这个对外监听器地址而 Spark 流任务同时写了localhost:9092,broker:29092其中broker:29092是容器内部监听器供 Spark 集群内的 executor 访问。核心PySpark Structured Streaming 作业streaming.py 是整套示例的核心它完整演示了 Structured Streaming 的「读流 → 解析 → 转换 → 多路输出」范式。第一步readStream 读取 Kafkadef 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_streamformat(kafka)依赖spark-sql-kafka-0-10连接器。读取出的原始 DataFrame 中key与value是binary类型还附带topic、partition、offset、timestamp等元数据列该 Schema 在 streaming-notebook.ipynb 的df_kafka_raw.printSchema()中有完整输出。第二步解析二进制消息为结构化列Kafka 消息只是逗号分隔的字符串需要按RIDE_SCHEMA解析def parse_ride_from_kafka_message(df, schema): assert 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], , ) # 按 Schema 逐个展开成顶层列并做类型转换 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])解析逻辑分三步先把 binary 的 key/value 转成 string再用F.split按, 切分最后依据RIDE_SCHEMA的字段顺序用col.getItem(idx).cast(field.dataType)逐列取出并转换类型。函数开头的assert df.isStreaming is True是一个很好的防御性检查——确保传入的确实是流式 DataFrame。解析后的 Schema 即为RIDE_SCHEMA的 7 列结构notebook 中有df_rides.printSchema()的完整输出佐证。第三步三种 Sink 与聚合操作脚本封装了三种输出方式与两类聚合Console Sink调试def sink_console(df, output_mode: str complete, processing_time: str 5 seconds): write_query df.writeStream \ .outputMode(output_mode) \ .trigger(processingTimeprocessing_time) \ .format(console) \ .option(truncate, False) \ .start() return write_querytrigger(processingTime5 seconds)采用固定间隔微批模式每 5 秒处理一批outputMode对非聚合流通常用append对聚合流用complete。Memory Sink交互式调试def sink_memory(df, query_name, query_template): query_df df.writeStream \ .queryName(query_name) \ .format(memory) \ .start() query_str query_template.format(table_namequery_name) query_results spark.sql(query_str) return query_results, query_dfMemory Sink 把结果写入内存中的临时表之后可以像普通表一样用spark.sql查询。notebook 中以vendor_id_counts为查询名执行select count(distinct(vendor_id)) from vendor_id_counts返回了结果2验证了内存表查询可用。Kafka Sink回写def sink_kafka(df, topic): write_query df.writeStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092,broker:29092) \ .outputMode(complete) \ .option(topic, topic) \ .option(checkpointLocation, checkpoint) \ .start() return write_queryKafka Sink 要求流 DataFrame 至少包含名为value的列若有 key 则还需key列因此先要经过prepare_df_to_kafka_sink加工def prepare_df_to_kafka_sink(df, value_columns, key_columnNone): 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])该函数用concat_ws(, , ...)把指定列重新拼回逗号分隔字符串作为value并把key_column此处为vendor_id重命名并转成 string 作为key。两类聚合算子def op_groupby(df, column_names): df_aggregation df.groupBy(column_names).count() return df_aggregation def op_windowed_groupby(df, window_duration, slide_duration): df_windowed_aggregation df.groupBy( F.window(timeColumndf.tpep_pickup_datetime, windowDurationwindow_duration, slideDurationslide_duration), df.vendor_id ).count() return df_windowed_aggregationop_groupby按vendor_id做普通分组计数输出到 Console Sinkcomplete模式op_windowed_groupby基于tpep_pickup_datetime时间列做滑动窗口聚合窗口 10 分钟、每 5 分钟滑动一次按窗口与vendor_id计数结果经过prepare_df_to_kafka_sink后回写到vendor_counts_windowed主题。主流程串联spark SparkSession.builder.appName(streaming-examples).getOrCreate() spark.sparkContext.setLogLevel(WARN) df_consume_stream read_from_kafka(consume_topicCONSUME_TOPIC_RIDES_CSV) df_rides parse_ride_from_kafka_message(df_consume_stream, RIDE_SCHEMA) sink_console(df_rides, output_modeappend) df_trip_count_by_vendor_id op_groupby(df_rides, [vendor_id]) df_trip_count_by_pickup_date_vendor_id op_windowed_groupby( df_rides, window_duration10 minutes, slide_duration5 minutes) sink_console(df_trip_count_by_vendor_id) 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) spark.streams.awaitAnyTermination()一个流式 DataFrame 可以被多个 sink 同时消费原始行程数据进 Consoleappend 模式按 vendor 计数进 Consolecomplete 模式窗口计数进 Kafkacomplete 模式。最后awaitAnyTermination()让主线程阻塞持续等待流查询终止。提交作业spark-submit.sh脚本无法直接python3 streaming.py运行因为 Kafka 连接器 JAR 不在 Spark 默认类路径中。spark-submit.sh 的作用就是在提交时通过--packages自动解析并下载所需依赖./spark-submit.sh streaming.py脚本内部执行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.3.1,org.apache.spark:spark-avro_2.12:3.3.1,org.apache.spark:spark-streaming-kafka-0-10_2.12:3.3.1 \ $PYTHON_JOB脚本还支持可选参数指定 executor 内存默认1G./spark-submit.sh streaming.py 2G提交参数含义参数值说明--masterspark://localhost:7077连接到 Docker 中 Spark Master 的 Standalone 集群地址--num-executors2使用两个 executor对应两个 spark-worker 容器--executor-memory1G默认每个 executor 的内存可按需传512M、2G等--executor-cores1每个 executor 的 CPU 核数--packages三个 3.3.1 版 JAR提供 Kafka 数据源/接收器与 Avro 支持其中spark-sql-kafka-0-10_2.12:3.3.1是 Structured Streaming 读写 Kafka 的必需连接器spark-avro用于 Avro 序列化场景本示例虽用 CSV 字符串但依赖已一并引入备用。在 notebook 中运行时则通过环境变量注入同样的依赖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依赖解析时 Ivy 会从 Maven Central 拉取这些工件及其传递依赖如kafka-clients、snappy-java等首次运行耗时较长属正常现象。需要留意的是这些版本号3.3.1、Kafka 7.2.0、Spark 3.3.1以当前仓库实际配置为准若本地 Spark 版本不同需相应调整。调试与排错要点从 streaming-notebook.ipynb 的运行输出中可以总结出几类典型问题Broker 不可达日志反复出现Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.说明 Kafka 容器未启动或监听器配置不匹配应先执行docker compose up -d并确认docker network ls中存在kafka-spark-network。startingOffsets 语义notebook 注释明确提示startingOffsets默认值为latest本示例显式设置为earliest以从头消费如果 Kafka 中没有历史数据且使用latest流会一直等待新消息。临时 checkpoint 告警Temporary checkpoint location created...属于正常行为未显式指定 checkpoint 时 Spark 会创建临时目录生产环境应像本示例那样用.option(checkpointLocation, checkpoint)固定位置以支持失败恢复。验证链路先启动consumer.py确认能收到 producer 发出的消息再提交streaming.py可将数据链路问题与 Spark 作业问题分开定位。小结本示例以 NYC 出租车行程数据为载体串起了流式处理最核心的四个环节Docker 编排 Kafka/Spark 环境、Python 生产者写入rides_csv主题、Structured Streaming 按RIDE_SCHEMA解析二进制消息并做分组与滑动窗口聚合、最后经 Console/Memory/Kafka 三种 Sink 输出结果。同一目录下的 streaming-notebook.ipynb 提供了逐 Cell 的交互式版本包含原始 Kafka Schema、编码后 Schema、结构化 Schema 及 Memory Sink 查询的真实输出可作为理解本文每一步的对照材料如需在笔记本环境中直接实验docker README 中的 JupyterLab 服务端口 8888与 Spark UI端口 8080均可配合使用。【免费下载链接】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),仅供参考
返回列表