【问题标题】:Spark saveAsTable throws NoSuchTableExceptionSpark saveAsTable 抛出 NoSuchTableException
【发布时间】:2020-04-02 12:11:52
【问题描述】:

我正在使用 pyspark 的 (Spark 2.3.2) saveAsTable 如下:

df.write.format("parquet") \
  .sortBy("id") \
  .bucketBy(50, "some_column") \
  .option("path", "test_table.parquet") \
  .saveAsTable("test_table", mode="overwrite")

在表已经存在的情况下(因此模式“覆盖”),这会导致NoSuchTableException

org.apache.spark.sql.catalyst.analysis.NoSuchTableException: Table or view 'test_table' not found in database 'test_database';
at org.apache.spark.sql.catalyst.catalog.SessionCatalog.requireTableExists(SessionCatalog.scala:184)
at org.apache.spark.sql.catalyst.catalog.SessionCatalog.listPartitionsByFilter(SessionCatalog.scala:927)
at org.apache.spark.sql.execution.datasources.CatalogFileIndex.filterPartitions(CatalogFileIndex.scala:73)
at org.apache.spark.sql.execution.datasources.CatalogFileIndex.listFiles(CatalogFileIndex.scala:59)
at org.apache.spark.sql.execution.FileSourceScanExec.org$apache$spark$sql$execution$FileSourceScanExec$$selectedPartitions$lzycompute(DataSourceScanExec.scala:189)
at org.apache.spark.sql.execution.FileSourceScanExec.org$apache$spark$sql$execution$FileSourceScanExec$$selectedPartitions(DataSourceScanExec.scala:186)
at org.apache.spark.sql.execution.FileSourceScanExec.inputRDD$lzycompute(DataSourceScanExec.scala:308)
at org.apache.spark.sql.execution.FileSourceScanExec.inputRDD(DataSourceScanExec.scala:295)
at org.apache.spark.sql.execution.FileSourceScanExec.inputRDDs(DataSourceScanExec.scala:315)
at org.apache.spark.sql.execution.WholeStageCodegenExec.doExecute(WholeStageCodegenExec.scala:605)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:131)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:127)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeQuery$1.apply(SparkPlan.scala:155)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:127)
at org.apache.spark.sql.execution.columnar.InMemoryRelation.buildBuffers(InMemoryRelation.scala:107)
at org.apache.spark.sql.execution.columnar.InMemoryRelation.<init>(InMemoryRelation.scala:102)
at org.apache.spark.sql.execution.columnar.InMemoryRelation$.apply(InMemoryRelation.scala:43)
at org.apache.spark.sql.execution.CacheManager.org$apache$spark$sql$execution$CacheManager$$recacheByCondition(CacheManager.scala:145)
at org.apache.spark.sql.execution.CacheManager$$anonfun$recacheByPath$1.apply$mcV$sp(CacheManager.scala:201)
at org.apache.spark.sql.execution.CacheManager$$anonfun$recacheByPath$1.apply(CacheManager.scala:194)
at org.apache.spark.sql.execution.CacheManager$$anonfun$recacheByPath$1.apply(CacheManager.scala:194)
at org.apache.spark.sql.execution.CacheManager.writeLock(CacheManager.scala:67)
at org.apache.spark.sql.execution.CacheManager.recacheByPath(CacheManager.scala:194)
at org.apache.spark.sql.internal.CatalogImpl.refreshByPath(CatalogImpl.scala:508)
at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:174)
at org.apache.spark.sql.execution.datasources.DataSource.writeAndRead(DataSource.scala:532)
at org.apache.spark.sql.execution.command.CreateDataSourceTableAsSelectCommand.saveDataIntoTable(createDataSourceTables.scala:216)
at org.apache.spark.sql.execution.command.CreateDataSourceTableAsSelectCommand.run(createDataSourceTables.scala:176)
at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult$lzycompute(commands.scala:104)
at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult(commands.scala:102)
at org.apache.spark.sql.execution.command.DataWritingCommandExec.doExecute(commands.scala:122)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:131)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:127)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeQuery$1.apply(SparkPlan.scala:155)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:127)
at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:80)
at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:80)
at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:656)
at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:656)
at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:77)
at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:656)
at org.apache.spark.sql.DataFrameWriter.createTable(DataFrameWriter.scala:458)
at org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:433)
at org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:393)
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:357)
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.GatewayConnection.run(GatewayConnection.java:238)
at java.lang.Thread.run(Thread.java:748)

看起来现有表已成功获得dropped,但following attempt to create the new table 似乎要求该表存在(参见堆栈跟踪的第二行)。这是一个错误还是我错过了什么?

【问题讨论】:

  • 我在覆盖 parquet 文件时遇到了类似的问题。会关注这个话题。

标签: apache-spark pyspark apache-spark-sql pyspark-sql


【解决方案1】:

它似乎在 Spark 2.4 中运行良好,主要是尝试创建一个示例数据帧,然后通过 spark 将其写入 Hive。

from pyspark.sql import Row
l = [('Ankit',25),('Jalfaizy',22),('Magesh',20),('Bala',26)]
rdd = sc.parallelize(l)
people = rdd.map(lambda x: Row(name=x[0], age=int(x[1])))
schemaPeople = spark.createDataFrame(people)
schemaPeople.write.format("parquet").saveAsTable("test_table_spark", mode="overwrite")

写入成功后,检查 Hive 表,然后修改数据帧,同样saveAsTable数据被新数据帧覆盖

l = [('Ankit',25),('Jalfaizy',22),('Suresh',20),('Bala',26)]

您能否在您的 spark shell 中尝试相同的方法,看看是否可行...

尝试在外部 Hive 表中执行相同操作

>>> schemaPeople.show()
+---+--------+
|age|    name|
+---+--------+
| 25|   Ankit|
| 22|Jalfaizy|
| 20|  Suresh|
| 26|    Bala|
+---+--------+

>>> spark.sql("SELECT * FROM EXT_Table_Test").show()
+---+--------+
|age|    name|
+---+--------+
| 25|   Ankit|
| 22|Jalfaizy|
| 20|  Magesh|
| 26|    Bala|
+---+--------+

>>> schemaPeople.write.format("parquet") \
...   .option("path", "hdfs://path/tables/EXT_Table_Test") \
...   .saveAsTable("test_table", mode="overwrite")

再次读取更新的表导致如下错误

原因:java.io.FileNotFoundException:文件不存在: hdfs:///tables/EXT_Table_Test/000000_0 有可能 基础文件已更新。您可以显式地使 通过在 SQL 中运行“REFRESH TABLE tableName”命令在 Spark 中缓存或 通过重新创建所涉及的 Dataset/DataFrame。

>>> spark.sql("SELECT * FROM EXT_Table_Test").show()

执行 REFRESH TABLE 后,从 Spark 读取成功,但在执行刷新之前,我能够在 HIVE shell 中看到更新的数据。

>>> spark.sql("REFRESH TABLE EXT_Table_Test")
DataFrame[]
>>> spark.sql("SELECT * FROM EXT_Table_Test").show()
+---+--------+
|age|    name|
+---+--------+
| 25|   Ankit|
| 22|Jalfaizy|
| 20|  Suresh|
| 26|    Bala|
+---+--------+

【讨论】:

  • 嗨@Joby,感谢您的意见。请注意,您的回答不考虑外部表 (.option("path", "/some/folder/file.parquet")),我认为这是问题的一部分。此外,如果以前不存在表,则写入表也不成问题。如果该表已存在,则似乎存在问题,尽管它确实被删除,但创建新表似乎存在问题。
  • @MartinStuder 为外部表更新了相同的内容,请检查是否有帮助
【解决方案2】:

在 spark 2.4 中创建覆盖失败的表。要解决此问题,请设置以下属性。

将标志 spark.sql.legacy.allowCreatingManagedTableUsingNonemptyLocation 设置为 true。

对于 pyspark,使用以下命令:

spark.conf.set("spark.sql.legacy.allowCreatingManagedTableUsingNonemptyLocation","true")

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-08-08
    • 1970-01-01
    • 2017-07-17
    • 2016-11-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多