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

