Spark Structured Streaming写入CSV失败:已设水位线仍报错
修复Spark Structured Streaming聚合时的水位线错误
环境信息
- Spark 3.4.1
- PySpark 3.4.1
- Python 3.11
问题场景
需要将结构化流DataFrame写入CSV文件,logsDF的Schema如下:
root |-- timestampLog: string (nullable = true) |-- status: integer (nullable = true) |-- timestampProcessing: timestamp (nullable = false)
已在代码中指定水位线,但运行时仍抛出“Append输出模式不支持无水位线的流聚合”的错误,代码如下:
statusCountsDF = logsDF \ .withWatermark("timestampProcessing", "10 minutes") \ .groupBy( window(logsDF.timestampProcessing, "10 minutes"), logsDF.status ).count() query = (statusCountsDF.writeStream .outputMode("append") .format("csv") .option("path", "logs/result") .option("header", True) .option("checkpointLocation", "logs/checkpoint") .queryName("counts") .start()) query.awaitTermination()
错误信息:
pyspark.errors.exceptions.captured.AnalysisException: Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark; Aggregate [window#15, status#3], [window#15 AS window#9, status#3, count(1) AS count#14L] +- Project [named_struct(start, knownnullable(precisetimestampconversion(((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - CASE WHEN (((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) < cast(0 as bigint)) THEN (((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) + 600000000) ELSE ((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) END) - 0), LongType, TimestampType)), end, knownnullable(precisetimestampconversion((((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - CASE WHEN (((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) < cast(0 as bigint)) THEN (((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) + 600000000) ELSE ((precisetimestampconversion(timestampProcessing#4, TimestampType, LongType) - 0) % 600000000) END) - 0) + 600000000), LongType, TimestampType))) AS window#15, timestampLog#2, status#3, timestampProcessing#4-T600000ms] +- Filter isnotnull(timestampProcessing#4) +- EventTimeWatermark timestampProcessing#4: timestamp, 10 minutes +- Project [regexp_extract(value#0, \[(\d{2}/\w{3}/\d{4}:\d{2}:\d{2}:\d{2} \+\d{4})], 1) AS timestampLog#2, cast(regexp_extract(value#0, \s(\d{3})\s, 1) as int) AS status#3, 2023-07-05 20:39:13.000826 AS timestampProcessing#4] +- StreamingRelation DataSource(org.apache.spark.sql.SparkSession@12625ab6,text,List(),None,List(),None,Map(wholetext -> false, path -> logs/),None), FileSource[logs/], [value#0]
修复方案
问题核心是聚合时引用了原始DataFrame的列,而非应用水位线后的DataFrame列。groupBy中使用logsDF.timestampProcessing会绕过水位线的处理逻辑,导致Spark无法识别水位线与聚合窗口的关联。
修改代码时,将groupBy中的logsDF.timestampProcessing改为直接引用列名,窗口函数基于应用水位线后的DataFrame:
from pyspark.sql.functions import window statusCountsDF = logsDF \ .withWatermark("timestampProcessing", "10 minutes") \ .groupBy( window("timestampProcessing", "10 minutes"), # 直接使用列名,而非logsDF.xxx "status" ).count() # 后续写入流代码保持不变 query = (statusCountsDF.writeStream .outputMode("append") .format("csv") .option("path", "logs/result") .option("header", True) .option("checkpointLocation", "logs/checkpoint") .queryName("counts") .start()) query.awaitTermination()
原理说明
调用withWatermark后会返回一个包含水位线元数据的新DataFrame。如果聚合时仍引用原始logsDF的列,Spark会判定聚合操作基于未应用水位线的原始流数据,因此触发“无水位线不支持Append模式”的错误。直接使用列名(或新DataFrame的列)能让Spark正确关联水位线与聚合窗口,从而允许使用Append输出模式。
内容的提问来源于stack exchange,提问作者jpseng
相关产品推荐
相关产品推荐

