
文章目录一、课前导读二、学习目标三、核心理论知识点四、原理通俗讲解4.1 微批次架构以时间换吞吐4.2 DStream与RDD的关系4.3 Structured Streaming的无限表模型五、重点概念拆解5.1 DStream API常用操作5.2 Structured Streaming核心概念5.3 窗口与水印六、易错点避坑6.1 混淆流处理与批处理的Checkpoint6.2 未设置水印导致状态无限增长6.3 输出模式选择错误6.4 DStream的foreachRDD中创建资源6.5 Kafka消费偏移量管理七、完整实战案例7.1 环境准备八、代码逐行解析8.1 DStream部分8.2 Structured Streaming部分九、业务场景落地应用9.1 场景一实时用户行为分析9.2 场景二实时风控9.3 场景三实时ETL9.4 场景四监控告警十、常见报错排查10.1 StreamingContext has already been started10.2 java.net.ConnectException: Connection refused Socket10.3 The output mode append is not allowed when there are streaming aggregations...10.4 Kafka消费时无法反序列化消息十一、本节课知识点总结Spark Streaming vs Structured Streaming关键参数十二、课后思考作业作业一理论理解题作业二代码实践题作业三场景应用题作业四拓展研究《20节课 PySpark 从入门到精通》系列课程导航一、课前导读在前面的课程中我们处理的都是静态的批量数据Batch Data从HDFS读取文件进行ETL输出结果任务运行完即结束。但在真实的企业场景中数据是持续不断地产生的——用户点击、订单支付、传感器读数、日志输出……如果你依然采用批处理的方式例如每小时运行一次那么数据的价值将大打折扣因为业务决策需要的是秒级甚至毫秒级的洞察。这就是**流式计算Streaming Computing的价值所在。Spark Streaming是Spark生态中最早用于实时流处理的组件它采用微批次Micro-Batch**的架构将连续的数据流切分成小批量然后复用Spark Core的批处理引擎进行运算。虽然严格来说它不是纯粹的“实时”延迟通常在秒级但它完美地结合了Spark的容错性和一致性是生产环境中最成熟的流处理方案之一。除了传统的Spark Streaming基于DStreamSpark 2.0后引入了Structured Streaming提供更简洁的API、事件时间处理、端到端精确一次Exactly-Once语义已经成为构建实时数据管道的新标准。本课程将同时讲解两种流处理API重点放在Structured Streaming上因为它代表了Spark流计算的未来。学完这节课你将能够理解流计算的架构原理并动手开发基于PySpark的实时流处理程序如实时日志分析、实时PV/UV统计、窗口聚合、与Kafka集成等。二、学习目标完成本节课的学习后你将能够理解Spark Streaming架构微批次原理、DStream模型、与批处理的关系掌握DStream API从Kafka/Socket创建流使用转换和输出操作深入理解Structured Streaming无限表模型、事件时间处理、水印机制实现窗口操作滚动窗口、滑动窗口处理延迟数据集成Kafka从Kafka消费数据处理并写回Kafka或外部存储管理状态使用groupBy和mapGroupsWithState维护流状态处理容错和一致性了解Checkpoint、Exactly-Once语义开发实战完成一个实时PV/UV统计系统三、核心理论知识点知识点说明微批次Spark Streaming将流数据切分成固定时间间隔的小批次如5秒每个批次作为RDD处理DStream离散化流表示连续的数据流内部是一系列RDD接收器Receiver传统Spark Streaming中从数据源接收数据的组件有可靠性之分Structured StreamingSpark 2.0流引擎基于Catalyst优化器支持事件时间和端到端Exactly-Once无限表Structured Streaming将流数据视为无限追加的表窗口操作基于事件时间的滚动窗口tumbling或滑动窗口sliding聚合水印Watermark处理延迟数据的机制允许一定的迟到时间决定何时输出最终结果状态存储维护流计算中的中间状态如聚合中的累计值可持久化到内存/磁盘Checkpoint保存流处理进度和状态用于故障恢复四、原理通俗讲解4.1 微批次架构以时间换吞吐想象你是一个工厂的质检员。如果每个零件一出来你就立即检测你可能会手忙脚乱效率很低。相反你用一个篮子接零件每隔5秒钟把篮子里积攒的一批零件一起检测。虽然每个零件的等待时间多了几秒但你的整体检测效率大幅提高。这就是微批次的精髓——用少量的延迟换取高吞吐和批处理引擎的复用。Spark Streaming的Driver会按照指定的间隔如5秒生成一个Jobs每个Job处理这个间隔内收集到的数据。数据在内部被切分成RDD的分区分布在Executor上并行处理。4.2 DStream与RDD的关系DStreamDiscretized Stream是Spark Streaming的核心抽象它本质上是一个连续的RDD序列每个批次一个RDD。你可以在DStream上应用map、filter、reduceByKey等转化操作这些操作会作用于底层的每个RDD并产生新的DStream。4.3 Structured Streaming的无限表模型Structured Streaming将流数据看作一个持续追加的无限表。每当有新数据到达Spark会“增量地”查询该表并将结果输出到“结果表”可以是内存、控制台、Kafka等。你可以用熟悉的DataFrame API或SQL来写流处理逻辑这与批处理几乎没有区别大大降低了学习门槛。当执行窗口聚合时Structured Streaming会根据事件时间和水印决定何时输出窗口的结果。五、重点概念拆解5.1 DStream API常用操作创建流frompyspark.streamingimportStreamingContext sscStreamingContext(sc,batchDuration2)# 2秒微批次# Socket文本流linesssc.socketTextStream(localhost,9999)# 或从Kafka (需引入依赖)转换map、flatMap、filter、reduceByKey、updateStateByKey等。输出print()打印前几个元素。saveAsTextFiles(prefix)保存到文件。foreachRDD自定义输出如写入数据库。启动和停止ssc.start()ssc.awaitTermination()5.2 Structured Streaming核心概念创建流# 读取Kafkadfspark.readStream.format(kafka)\.option(kafka.bootstrap.servers,host:9092)\.option(subscribe,topic)\.load()转换使用DataFrame API或SQL与批处理几乎相同。窗口聚合frompyspark.sql.functionsimportwindow,col df.groupBy(window(timestamp,10 seconds,5 seconds)).count()输出模式append只输出新行默认适用于无聚合。update输出有更新的行。complete输出全量结果表适用于聚合。输出querydf.writeStream \.outputMode(append)\.format(console)\.start()query.awaitTermination()5.3 窗口与水印滚动窗口Tumbling固定大小不重叠如每5分钟统计一次。滑动窗口Sliding有重叠需要指定窗口长度和滑动步长如窗口10分钟滑动5分钟。水印定义允许的延迟时间如withWatermark(timestamp, 10 minutes)通知引擎对于超过10分钟的数据不再更新窗口结果。六、易错点避坑6.1 混淆流处理与批处理的CheckpointSpark Streaming的Checkpoint用于保存元数据和状态如updateStateByKey防止Driver重启后丢失。必须设置ssc.checkpoint(hdfs://path)。6.2 未设置水印导致状态无限增长在Structured Streaming的聚合中如果不设置水印状态会累积所有历史数据导致内存溢出。必须设置水印来清理旧状态。6.3 输出模式选择错误append模式要求聚合结果不会更新如水印后如果聚合结果可能更新如没有水印的聚合会抛出异常。6.4 DStream的foreachRDD中创建资源应在foreachRDD内部创建连接如数据库连接并确保正确关闭避免连接泄漏。6.5 Kafka消费偏移量管理使用Structured Streaming时偏移量由Checkpoint自动管理。如果应用升级需注意Checkpoint兼容性或重置偏移量。七、完整实战案例本案例将实现两个示例1使用传统Spark StreamingDStream进行socket实时单词计数2使用Structured Streaming模拟实时日志分析计算每分钟的PV、UV并输出到控制台。7.1 环境准备需要安装pyspark无需额外依赖。测试时可用nc -lk 9999发送数据。# streaming_demo.py # 功能Spark Streaming 与 Structured Streaming 实战# 包含DStream WordCount、Structured Streaming PV/UV统计、窗口与水印frompyspark.sqlimportSparkSessionfrompyspark.streamingimportStreamingContextfrompyspark.sql.functionsimport*frompyspark.sql.typesimport*importtime# 1. 传统 DStream 示例WordCount from socket defdstream_wordcount():print(*80)print(传统 Spark Streaming (DStream) WordCount 示例)print(请在另一个终端运行: nc -lk 9999)print(*80)sparkSparkSession.builder \.appName(DStreamWordCount)\.master(local[2])\.getOrCreate()scspark.sparkContext sscStreamingContext(sc,batchDuration5)# 5秒一个批次# 创建DStreamlinesssc.socketTextStream(localhost,9999)# 处理逻辑wordslines.flatMap(lambdaline:line.split( ))pairswords.map(lambdaword:(word,1))word_countspairs.reduceByKey(lambdaa,b:ab)# 打印结果word_counts.pprint()# 启动ssc.start()print(监听中... 10秒后自动停止)time.sleep(30)ssc.stop(stopSparkContextTrue,stopGraceFullyTrue)print(DStream示例结束)# 注意实际运行时需手动发送数据这里模拟调用# dstream_wordcount()# 2. Structured Streaming 示例实时PV/UV统计 defstructured_streaming_pvuv():print(\n*80)print(Structured Streaming 实时PV/UV统计)print(模拟生成用户访问日志计算每分钟PV和UV)print(*80)sparkSparkSession.builder \.appName(StructuredStreamingPVUV)\.master(local[2])\.config(spark.sql.shuffle.partitions,2)\.getOrCreate()scspark.sparkContext sc.setLogLevel(WARN)# 生成模拟数据流每秒生成随机访问记录frompyspark.sql.functionsimportstruct,to_jsonimportrandomfromdatetimeimportdatetime,timedelta# 创建一个自定义的数据源生成器使用RateStream或MemoryStream这里用Rate模拟# 方式1: 使用内置的rate source每秒生成数据# 但rate source不包含我们需要的字段我们使用for each batch方式生成模拟# 更简单使用readStream.format(rate)生成时间戳然后用withColumn添加用户ID等# 生成模拟访问日志的流数据# 方式使用MemoryStream测试用或自定义DataFrame流# 为了演示我们创建一个持续生成的DataFrame通过foreachBatch不断追加# 但Structured Streaming需要真实的数据源这里我们用socket模拟实际场景# 实际生产应连接Kafka这里为了方便用socket模拟用户访问日志# 启动一个socket server发送JSON格式的日志# 我们使用内置的rate source 生成随机字段作为替代print(启动模拟访问日志生成器...)# 使用rate source产生时间戳每秒生成10条rate_dfspark.readStream.format(rate)\.option(rowsPerSecond,10)\.option(numPartitions,2)\.load()# 为每条数据添加随机用户ID1-1000和URLlogs_dfrate_df.withColumn(user_id,(rand()*1000).cast(int)1)\.withColumn(url,concat(lit(/page/),(rand()*10).cast(int)))\.withColumn(timestamp,col(timestamp))\.select(timestamp,user_id,url)# 定义窗口1分钟滚动窗口windowed_countslogs_df \.withWatermark(timestamp,10 seconds)\.groupBy(window(timestamp,1 minute),url)\.agg(count(*).alias(pv),approx_count_distinct(user_id).alias(uv))\.select(window.start,window.end,url,pv,uv)# 输出到控制台querywindowed_counts.writeStream \.outputMode(append)\.format(console)\.option(truncate,false)\.trigger(processingTime5 seconds)\.start()print(开始实时统计将每分钟输出一次窗口结果运行30秒...)time.sleep(30)query.stop()# 另一种方式从Kafka读取真实场景print(\n演示从Kafka读取的配置注释代码)print( kafka_df spark.readStream.format(kafka) \\ .option(kafka.bootstrap.servers, localhost:9092) \\ .option(subscribe, user-logs) \\ .load() \\ .selectExpr(CAST(value AS STRING) as json) \\ .select(from_json(json, schema).alias(data)) \\ .select(data.*) )spark.stop()print(Structured Streaming示例结束)# 3. 高级状态操作与窗口聚合 defstructured_streaming_aggregate_state():print(\n*80)print(高级状态操作 - 按用户聚合点击次数使用groupBy agg)print(*80)sparkSparkSession.builder \.appName(StructuredState)\.master(local[2])\.getOrCreate()# 模拟流数据frompyspark.sql.functionsimportrand,struct,to_json sourcespark.readStream.format(rate).option(rowsPerSecond,5).load()clickssource.withColumn(user_id,(rand()*100).cast(int))\.withColumn(event,lit(click))\.select(timestamp,user_id,event)# 按用户分组统计点击次数无窗口持续更新# 需要设置水印避免状态无限增长user_countsclicks.withWatermark(timestamp,1 hour)\.groupBy(user_id)\.count()\.select(user_id,count)# 输出更新模式queryuser_counts.writeStream \.outputMode(update)\.format(console)\.trigger(processingTime5 seconds)\.start()print(每个用户的累计点击次数运行20秒...)time.sleep(20)query.stop()spark.stop()# 4. 实际生产场景消费Kafka写入HDFS defproduction_kafka_to_hdfs():print(\n*80)print(生产场景示例从Kafka消费日志清洗后写入HDFSParquet)print(*80)code # 定义JSON Schema from pyspark.sql.types import StructType, StructField, StringType, TimestampType schema StructType([ StructField(user_id, StringType()), StructField(action, StringType()), StructField(timestamp, TimestampType()) ]) # 读取Kafka df spark.readStream.format(kafka) \\ .option(kafka.bootstrap.servers, kafka:9092) \\ .option(subscribe, app-logs) \\ .load() \\ .selectExpr(CAST(value AS STRING) as json) \\ .select(from_json(json, schema).alias(data)) \\ .select(data.*) # 清洗、转换 df df.filter(col(action).isNotNull()) # 写入HDFS按日期分区 query df.writeStream \\ .outputMode(append) \\ .format(parquet) \\ .option(path, /warehouse/ods/logs) \\ .option(checkpointLocation, /checkpoints/logs) \\ .partitionBy(year, month, day) \\ .trigger(processingTime1 minute) \\ .start() print(code)# 主函数 if__name____main__:# 依次运行示例需要取消注释并手动测试# dstream_wordcount()structured_streaming_pvuv()# structured_streaming_aggregate_state()production_kafka_to_hdfs()八、代码逐行解析8.1 DStream部分sscStreamingContext(sc,batchDuration5)linesssc.socketTextStream(localhost,9999)word_countslines.flatMap(...).map(...).reduceByKey(...)word_counts.pprint()ssc.start()StreamingContext是入口需要SparkContext和批次间隔。socketTextStream创建从socket接收文本的DStream。转换算子与RDD类似但作用在DStream上。pprint()打印前10个元素。start()启动流引擎awaitTermination阻塞等待。8.2 Structured Streaming部分rate_dfspark.readStream.format(rate).option(rowsPerSecond,10).load()rate是内置数据源用于测试每秒生成指定条数数据。withWatermark(timestamp, 10 seconds)定义了水印允许延迟10秒。groupBy(window(timestamp, 1 minute), url)按1分钟滚动窗口和URL分组。输出模式append因为窗口结果确定后不会更新。trigger(processingTime5 seconds)每5秒触发一次计算。九、业务场景落地应用9.1 场景一实时用户行为分析埋点日志进入KafkaStructured Streaming消费后按用户ID聚合计算每个用户最近1小时的点击序列用于实时推荐。9.2 场景二实时风控支付流水流通过窗口聚合如1分钟内同一用户下单金额累计超过阈值触发风控告警。9.3 场景三实时ETL将Canal捕获的MySQL binlogJSON格式实时同步到Hudi/Delta Lake构建实时数据湖。9.4 场景四监控告警服务器日志流解析错误关键词若5分钟内错误数超过10次则发送钉钉告警。十、常见报错排查10.1StreamingContext has already been started原因重复调用ssc.start()或之前未停止。解决确保只启动一次用ssc.awaitTermination()保持运行。10.2java.net.ConnectException: Connection refusedSocket原因没有开启对应的socket server。解决用nc -lk 9999启动。10.3The output mode append is not allowed when there are streaming aggregations...原因聚合操作的结果可能更新但使用了append模式。解决改为update或complete模式或设置水印使结果最终确定。10.4 Kafka消费时无法反序列化消息原因value是字节数组需显式转换为字符串。解决selectExpr(CAST(value AS STRING))。十一、本节课知识点总结Spark Streaming vs Structured Streaming特性DStreamStructured Streaming模型RDD序列无限表 DataFrameAPI复杂度较旧需理解RDDSQL/DataFrame简单事件时间支持弱原生支持水印窗口滑动窗口需指定滚动/滑动支持延迟数据一致性至少一次可配置精确一次精确一次性能高高且有Catalyst优化关键参数spark.sql.streaming.schemaInference是否推断Schema流中慎用。spark.sql.streaming.stateStore.providerClass状态存储实现。spark.sql.shuffle.partitions流聚合的shuffle分区数。十二、课后思考作业作业一理论理解题解释Spark Streaming微批次架构中批次间隔与窗口长度的关系。Structured Streaming的水印机制如何解决迟到的数据问题输出模式update和complete的区别是什么分别适用于哪些场景作业二代码实践题编写一个DStream程序从Kafka消费JSON格式的用户登录事件统计每5分钟每个IP的登录失败次数超过3次则输出告警。使用Structured Streaming实现一个实时热门商品统计数据流包含user_id, product_id, event_time计算每分钟点击量Top3的商品输出到控制台。模拟一个从socket读取数据每行格式时间戳用户金额使用Structured Streaming计算每10秒滚动窗口的用户消费总额并输出到内存中。作业三场景应用题某游戏公司需要实时分析用户充值流水每秒数千条。需求实时统计每1分钟每个区服的总充值金额和人数检测单用户单日累计充值超过1000元的立即输出风控事件将聚合结果写入MySQL数据库用于大屏展示请设计基于Spark Structured Streaming的方案包括Kafka topic设计、窗口大小、水印设置、输出模式、Checkpoint配置。作业四拓展研究调研Structured Streaming与Flink的对比分析各自的优劣和适用场景。学习Spark Structured Streaming的mapGroupsWithState和flatMapGroupsWithStateAPI实现一个会话窗口session window的统计。搭建Kafka Spark Streaming Elasticsearch的实时日志分析管道写出详细步骤。提交方式本次作业要求提交可运行的代码和运行日志以及理论题的解答。建议使用Kafka和Spark集群进行测试。扩展阅读Spark官方文档Streaming Programming GuideStructured Streaming Programming Guide《Stream Processing with Apache Spark》通过本节课的学习你已经掌握了PySpark实时流计算的核心技术能够开发低延迟的流处理应用。下一节课我们将学习PySpark性能全维度调优包括参数调优、SQL调优、数据倾斜根治等高级技巧助你成为Spark性能优化专家。我们下节课见《20节课 PySpark 从入门到精通》系列课程导航去订阅 感谢您耐心阅读到这里 如果本文对您有所启发欢迎 点赞 收藏 分享给更多需要的伙伴。️ 期待在评论区看到您的想法, 共同进步。 关注我持续获取更多干货内容 我们下篇文章见