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
相关产品推荐
相关产品推荐

