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

Dataproc 2.2环境下Spark对流式DataFrame执行聚合操作时抛出IOException的问题咨询

Dataproc 2.2环境下Spark对流式DataFrame执行聚合操作时抛出IOException的问题咨询

我最近在把一个运行在Dataproc 2.1镜像(Spark 3.3、Python 3.10)上的作业迁移到Dataproc 2.2镜像(Spark 3.5、Python 3.11)时遇到了问题,其中一个流式查询抛出了IO异常。我已经把问题简化成了最小复现示例:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F


spark = SparkSession.Builder().getOrCreate()

df =  (spark.readStream.format("rate")
        .option("rowsPerSecond", 4)
        .load()
    ).withWatermark(
        "timestamp", "1 seconds"
    ).withColumn(
        "window",
        F.window("timestamp", "10 seconds"),
    ).groupBy(F.col("window")).agg(F.count(F.expr("*")))


df.writeStream.format("console").queryName("sql_console").start()
spark.streams.awaitAnyTermination(60)
_ = [q.stop() for q in  spark.streams.active]

运行这段代码会触发如下异常:

25/03/17 11:37:04 WARN TaskSetManager: Lost task 7.0 in stage 1.0 (TID 7) (cluster-test-alexis-<redacted>.internal executor 2): java.io.IOException: mkdir of file:/tmp/temporary-42a491af-ba8c-412b-b06e-fb420879b92f/state/0/7 failed
    at org.apache.hadoop.fs.FileSystem.primitiveMkdir(FileSystem.java:1414)
    at org.apache.hadoop.fs.DelegateToFileSystem.mkdir(DelegateToFileSystem.java:185)
    at org.apache.hadoop.fs.FilterFs.mkdir(FilterFs.java:219)
    at org.apache.hadoop.fs.FileContext$4.next(FileContext.java:818)
    [...]

我发现如果去掉groupBy(...).agg(...)这部分逻辑,异常就会消失。初步猜测是Spark在执行聚合操作时需要创建临时文件,但创建失败了——不过暂时不确定是权限不足、磁盘空间不够还是其他原因导致的。


补充日志信息1

在查看作业日志时,我还在启动阶段发现了这条警告:

25/03/17 11:36:31 WARN ResolveWriteToStream: Temporary checkpoint location created which is deleted normally when the query didn't fail: /tmp/temporary-42a491af-ba8c-412b-b06e-fb420879b92f. If it's required to delete it under any circumstances, please set spark.sql.streaming.forceDeleteTempCheckpointLocation to true. Important to know deleting temp checkpoint folder is best effort.

根据这条信息,临时检查点目录应该已经成功创建了,那是不是说明问题不是出在权限上?毕竟如果连初始目录都创建不了的话,这条警告应该也不会出现?


临时解决方案

后来我找到了一个简单的临时 workaround:显式指定检查点位置,修改后的SparkSession初始化代码如下:

spark = SparkSession.Builder().config("spark.sql.streaming.checkpointLocation", "/tmp").getOrCreate()

另外我注意到,在HDFS的/tmp目录下并没有看到以temporary-开头的文件夹,只有以查询ID命名的文件夹。

现在我想请教:我该如何进一步排查这个问题的根本原因?有没有更完善的解决方案?

备注:内容来源于stack exchange,提问作者AlexisBRENON

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 19:03:10