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

为何PySpark作业执行最终checkpoint时失败?Spark3.5.6复现

Spark 3.5.6 大小写不敏感场景下Union与Checkpoint组合触发的异常

以下测试基于Spark 3.5.6:

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("reproduce_union_checkpoint_example")
    .config("spark.sql.caseSensitive","false")
    .getOrCreate()
)
spark.sparkContext.setCheckpointDir("/tmp/foo") 

df_1 = spark.createDataFrame([("a",)], schema=["foo"])
df_2 = df_1.where("Foo == 'a'").checkpoint() #注意!此处where条件使用了大写的Foo
df_2.printSchema() #此处显示列名为foo
df_union = df_1.unionByName(df_2)
df_union.explain(True)
df_union.show() #此步骤正常执行
df_union = df_union.checkpoint() #此步骤失败
df_union.show()

复现关键条件

  • 若df_2不基于df_1,作业执行成功。
  • 若df_2不与df_1执行union操作,查询执行成功。
  • 若省略任意一个checkpoint(),查询执行成功。
  • 若在where()中使用正确大小写的列名(写成"foo == 'a'"),查询执行成功。

执行输出

root
 |-- foo: string (nullable = true)

== Parsed Logical Plan ==
Union false, false
:- LogicalRDD [foo#0], false
+- Project [foo#2]
   +- LogicalRDD [foo#2], false

== Analyzed Logical Plan ==
foo: string
Union false, false
:- LogicalRDD [foo#0], false
+- Project [foo#2]
   +- LogicalRDD [foo#2], false

== Optimized Logical Plan ==
Union false, false
:- LogicalRDD [foo#0], false
+- LogicalRDD [foo#2], false

== Physical Plan ==
Union
:- *(1) Scan ExistingRDD[foo#0]
+- *(2) Scan ExistingRDD[foo#2]

+---+
|foo|
+---+
|  a|
|  a|
+---+

Traceback (most recent call last):
  File "/home/user/checkpoint_new_short.py", line 17, in <module>
    df_union = df_union.checkpoint()
  File "/usr/lib/spark3/python/lib/pyspark.zip/pyspark/sql/dataframe.py", line 1038, in checkpoint
  File "/usr/lib/spark3/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1323, in __call__
  File "/usr/lib/spark3/python/lib/pyspark.zip/pyspark/errors/exceptions/captured.py", line 179, in deco
  File "/usr/lib/spark3/python/lib/py4j-0.10.9.7-src.zip/py4j/protocol.py", line 328, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o129.checkpoint.
: java.util.NoSuchElementException: key not found: Foo#0
        at scala.collection.MapOps.default(Map.scala:274)
        at scala.collection.MapOps.default$(Map.scala:273)
        at org.apache.spark.sql.catalyst.expressions.AttributeMap.default(AttributeMap.scala:41)
        at scala.collection.MapOps.apply(Map.scala:176)
        at scala.collection.MapOps.apply$(Map.scala:175)
        at org.apache.spark.sql.catalyst.expressions.AttributeMap.apply(AttributeMap.scala:41)
        at org.apache.spark.sql.catalyst.plans.logical.Union$$anonfun$$nestedInanonfun$rewriteConstraints$1$1.applyOrElse(basicLogicalOperators.scala:515)
        at org.apache.spark.sql.catalyst.plans.logical.Union$$anonfun$$nestedInanonfun$rewriteConstraints$1$1.applyOrElse(basicLogicalOperators.scala:514)
        at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:461)
        at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(origin.scala:76)
        at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:461)
        at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$3(TreeNode.scala:466)
        at org.apache.spark.sql.catalyst.trees.UnaryLike.mapChildren(TreeNode.scala:1216)
        at org.apache.spark.sql.catalyst.trees.UnaryLike.mapChildren$(TreeNode.scala:1215)
        at org.apache.spark.sql.catalyst.expressions.UnaryExpression.mapChildren(Expression.scala:533)
        at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:466)
        at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:437)
        at org.apache.spark.sql.catalyst.trees.TreeNode.transform(TreeNode.scala:405)
        at org.apache.spark.sql.catalyst.plans.logical.Union.$anonfun$rewriteConstraints$1(basicLogicalOperators.scala:514)
        at org.apache.spark.sql.catalyst.expressions.ExpressionSet.$anonfun$map$1(ExpressionSet.scala:149)
        at org.apache.spark.sql.catalyst.expressions.ExpressionSet.$anonfun$map$1$adapted(ExpressionSet.scala:149)
        at scala.collection.IterableOnceOps.foreach(IterableOnce.scala:563)
        at scala.collection.IterableOnceOps.foreach$(IterableOnce.scala:561)
        at scala.collection.AbstractIterator.foreach(Iterator.scala:1293)
        at org.apache.spark.sql.catalyst.expressions.ExpressionSet.map(ExpressionSet.scala:149)
        at org.apache.spark.sql.catalyst.plans.logical.Union.rewriteConstraints(basicLogicalOperators.scala:514)
        at org.apache.spark.sql.catalyst.plans.logical.Union.$anonfun$validConstraints$2(basicLogicalOperators.scala:535)
        at scala.collection.immutable.Vector1.map(Vector.scala:1886)
        at scala.collection.immutable.Vector1.map(Vector.scala:375)
        at org.apache.spark.sql.catalyst.plans.logical.Union.validConstraints$lzycompute(basicLogicalOperators.scala:535)
        at org.apache.spark.sql.catalyst.plans.logical.Union.validConstraints(basicLogicalOperators.scala:533)
        at org.apache.spark.sql.catalyst.plans.logical.QueryPlanConstraints.constraints(QueryPlanConstraints.scala:34)
        at org.apache.spark.sql.catalyst.plans.logical.QueryPlanConstraints.constraints$(QueryPlanConstraints.scala:32)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.constraints$lzycompute(LogicalPlan.scala:32)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.constraints(LogicalPlan.scala:32)
        at org.apache.spark.sql.execution.LogicalRDD$.$anonfun$rewriteStatsAndConstraints$1(ExistingRDD.scala:224)
        at scala.Option.map(Option.scala:242)
        at org.apache.spark.sql.execution.LogicalRDD$.rewriteStatsAndConstraints(ExistingRDD.scala:222)
        at org.apache.spark.sql.execution.LogicalRDD$.fromDataset(ExistingRDD.scala:186)
        at org.apache.spark.sql.Dataset.$anonfun$checkpoint$1(Dataset.scala:741)
        at org.apache.spark.sql.Dataset.$anonfun$withAction$2(Dataset.scala:4324)
        at org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:546)
        at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:4322)
        at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:125)
        at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:201)
        at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:108)
        at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:900)
        at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66)
        at org.apache.spark.sql.Dataset.withAction(Dataset.scala:4322)
        at org.apache.spark.sql.Dataset.checkpoint(Dataset.scala:727)
        at org.apache.spark.sql.Dataset.checkpoint(Dataset.scala:690)
        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
        at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374)
        at py4j.Gateway.invoke(Gateway.java:282)
        at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
        at py4j.commands.CallCommand.execute(CallCommand.java:79)
        at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
        at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
        at java.lang.Thread.run(Thread.java:750)

内容的提问来源于stack exchange,提问作者Alexey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:37:02