You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Structured Streaming Parquet聚合后无输出问题求助

问题描述

我有一个存储在HDFS上的Parquet格式数据集,Schema如下:

IDNameRating
42"Book name 1""liked it"
57"Book name 2""really liked it"

我需要统计每本书的行数并将结果写入另一个Parquet文件。用控制台输出执行聚合时一切正常:

rows_count_query =\
    spark\
        .readStream\
        .schema(user_rating_schema)\
        .parquet(path=src_data_path)\
        .groupBy("Name")\
        .agg(fun.count(fun.col("Name")).alias("RowCount"))\
        .writeStream\
        .format("console")\
        .outputMode("complete")\
        .start()

输出正常显示统计结果。但尝试写入Parquet时,读取输出路径发现数据为空:

rows_count_query =\
    spark\
        .readStream\
        .schema(user_rating_schema)\
        .parquet(path=input_path)\
        .select("Name", fun.current_timestamp().alias("CurrentTime"))\
        .withWatermark("CurrentTime", delayThreshold="1 minute")\
        .groupBy("Name", "CurrentTime")\
        .agg(fun.count(fun.col("Name")).alias("RowCount"))\
        .writeStream\
        .format("parquet")\
        .option("path", output_path)\
        .option("checkpointLocation", "/tmp/checkpoint")\
        .outputMode("append")\
        .start()

读取输出的代码:

df = spark.read.parquet(output_path)
df.show()

输出为空。我添加了current_timestamp()作为水印所需的时间戳,以支持append输出模式的聚合。为何输出文件无内容?如何解决才能将结果写入Parquet?

原因分析与解决方案

核心问题

  1. 水印触发条件未达标:水印设置的1分钟延迟要求,只有当事件时间(此处为CurrentTime)超过当前最大事件时间+1分钟时,Spark才会输出聚合结果。但静态数据集的处理几乎瞬间完成,没有足够时间满足水印触发逻辑。
  2. 分组逻辑错误:按Name+CurrentTime分组,而CurrentTime是每条数据的实时处理时间,会导致同一本书的每条记录被分到不同组,聚合结果为单条计数1,且因水印未触发无法输出。

解决方法

方案一:改用complete输出模式(适合静态/小批量流数据)

直接复用控制台输出的聚合逻辑,用complete模式输出全量聚合结果,无需水印和额外时间戳:

rows_count_query =\
    spark\
        .readStream\
        .schema(user_rating_schema)\
        .parquet(path=input_path)\
        .groupBy("Name")\
        .agg(fun.count(fun.col("Name")).alias("RowCount"))\
        .writeStream\
        .format("parquet")\
        .option("path", output_path)\
        .option("checkpointLocation", "/tmp/checkpoint")\
        .outputMode("complete")\
        .start()

该模式会将全量聚合结果写入输出,适合需要完整统计的场景。

方案二:调整水印与分组逻辑(适合持续流数据)

若必须使用append模式,需修改分组逻辑并调整水印延迟:

rows_count_query =\
    spark\
        .readStream\
        .schema(user_rating_schema)\
        .parquet(path=input_path)\
        .select("Name", fun.current_timestamp().alias("CurrentTime"))\
        .withWatermark("CurrentTime", delayThreshold="0 seconds")\
        .groupBy("Name")\  # 仅按Name分组,移除CurrentTime
        .agg(fun.count(fun.col("Name")).alias("RowCount"))\
        .writeStream\
        .format("parquet")\
        .option("path", output_path)\
        .option("checkpointLocation", "/tmp/checkpoint")\
        .outputMode("append")\
        .start()

同时需等待流处理触发输出,比如调用rows_count_query.awaitTermination(10)设置10秒超时,让Spark完成数据写入。

额外注意事项

  • 读取输出前需确保流处理已完成或触发输出,避免在数据未写入时读取导致空结果。
  • 若之前查询有失败记录,需清理/tmp/checkpoint目录,避免残留状态影响新查询。

内容的提问来源于stack exchange,提问作者Георгий Гуминов

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 15:35:28