PySpark中spark.read()与spark.readStream()的区别是什么?
Spark.read() vs Spark.readStream():差异与适用场景
核心定位
spark.read():Spark批处理API,专门用于一次性读取静态数据集,处理完成后任务直接结束。spark.readStream():Spark结构化流API,用于持续读取实时数据流,会一直运行并处理新产生的数据,直到手动停止。
关键差异
- 处理模式
spark.read():全量读取数据,属于一次性批处理,作业执行完就终止。spark.readStream():增量读取新数据,属于持续流处理,作业长期运行。
- 数据源支持
- 两者都兼容Parquet、CSV、JSON、JDBC等常规数据源,但
spark.readStream()额外支持Kafka、Kinesis、目录监控这类能持续生成新数据的流数据源。
- 两者都兼容Parquet、CSV、JSON、JDBC等常规数据源,但
- 输出逻辑
spark.read()处理后直接输出到静态存储或数据库,一次完成写入。spark.readStream()需要配置输出模式(Append/Update/Complete),持续将增量结果输出,还能支持更新聚合表这类动态操作。
- 容错与状态
spark.read()无状态,失败重启会从头读取全量数据重新处理。spark.readStream()自带状态管理,支持Exactly-Once语义,重启后能从上次中断的位置继续处理,不会丢数据或重复处理。
适用场景
用spark.read()的情况
- 处理历史静态数据,比如分析过去一年的销售订单记录。
- 一次性数据清洗、转换任务,比如把CSV格式的历史日志转成Parquet格式存储。
- 不需要实时响应的离线数据分析,比如周度/月度报表生成。
用spark.readStream()的情况
- 实时数据监控,比如实时统计网站用户的PV/UV指标。
- 实时告警系统,比如监控服务器CPU、内存指标,异常时触发告警。
- 持续数据集成,比如从Kafka实时读取业务数据,同步到数据仓库或实时报表系统。
代码示例
批处理(spark.read())
# 读取CSV格式的历史订单数据 df = spark.read.csv("/dbfs/historical_orders.csv", header=True, inferSchema=True) # 统计各地区订单量 region_order_counts = df.groupBy("region").count() # 输出结果到Parquet文件 region_order_counts.write.parquet("/dbfs/region_order_results.parquet")
流处理(spark.readStream())
# 从Kafka读取实时订单流数据 stream_df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafka-cluster:9092") \ .option("subscribe", "real-time-orders-topic") \ .load() # 解析Kafka消息内容 parsed_stream = stream_df.selectExpr("CAST(value AS STRING) AS order_info") # 按10分钟窗口统计各地区实时订单量 windowed_stats = parsed_stream.groupBy( window("event_time", "10 minutes"), "region" ).count() # 以Append模式输出结果到控制台(生产环境可写入数据库或流存储) stream_query = windowed_stats.writeStream \ .outputMode("append") \ .format("console") \ .start() # 保持流作业运行 stream_query.awaitTermination()
内容的提问来源于stack exchange,提问作者Vinit Khandelwal
相关产品推荐
相关产品推荐

