【问题标题】:How to make concurrency work with dataframes writing into hive tables?如何使并发与写入配置单元表的数据帧一起工作?
【发布时间】:2019-05-26 04:54:19
【问题描述】:

我在 Spark 1.6 上有多个线程写入同一个 hive 表(使用 parquet 文件),当它们尝试同时写入时,在将写入文件的部分重命名为 HDFS 期间会提示错误。我正在寻找绕过这个已知 Spark 问题的解决方案。

class MyThread extends Runnable {
          def run {
          //some code
          myTable.write.format("parquet").mode("append")
                 .saveAsTable("hdfstable")
          //some code
          }
}
Executors.defaultThreadFactory().newThread(new MyThread).start()

我收到此错误:

org.apache.spark.SparkException: Job aborted.
    at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelation$$anonfun$run$1.apply$mcV$sp(InsertIntoHadoopFsRelation.scala:156)
    at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelation$$anonfun$run$1.apply(InsertIntoHadoopFsRelation.scala:108)
    at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelation$$anonfun$run$1.apply(InsertIntoHadoopFsRelation.scala:108)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:53)
    at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelation.run(InsertIntoHadoopFsRelation.scala:108)
    at org.apache.spark.sql.execution.ExecutedCommand.sideEffectResult$lzycompute(commands.scala:58)
    at org.apache.spark.sql.execution.ExecutedCommand.sideEffectResult(commands.scala:56)
    at org.apache.spark.sql.execution.ExecutedCommand.doExecute(commands.scala:70)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$5.apply(SparkPlan.scala:132)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$5.apply(SparkPlan.scala:130)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:150)
    at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:130)
    at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:55)
    at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:55)
    at org.apache.spark.sql.DataFrameWriter.insertInto(DataFrameWriter.scala:189)
    at org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:239)
    at org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:221)
    at fr.neolink.spark.streaming.StreamingNeo$.algo(StreamingNeo.scala:837)
    at fr.neolink.spark.streaming.StreamingNeo$$anonfun$main$3$$anonfun$apply$18$MyThread$1.run(StreamingNeo.scala:374)
    at java.lang.Thread.run(Thread.java:748)

引起:

java.io.IOException: Failed to rename 
FileStatus{path=hdfs://my_hdfs_master/user/hive/warehouse/MYDB.db/hdfstable/_temporary/0/task_201812281010_1770_m_000000/part-r-00000-9a70cbea-d105-4f50-ba1b-372f555906ce.gz.parquet; 
isDirectory=false; length=4608; replication=3; blocksize=134217728; modification_time=1545988247575; 
access_time=1545988247494; owner=owner; group=hive; permission=rw-r--r--; isSymlink=false} 
to hdfs://my_hdfs_master/user/hive/warehouse/MYDB.db/hdfstable/part-r-00000-9a70cbea-d105-4f50-ba1b-372f555906ce.gz.parquet

我在 jira 上发现了这个问题:https://issues.apache.org/jira/browse/SPARK-18626

有没有办法让写作部分线程安全?一个接一个地执行?

谢谢。

【问题讨论】:

  • 恐怕不行。我们已经在我们的生产环境中尝试过它,有时它会从 hive 目录中出现死锁。所以我建议你一个一个地执行,或者使用多个进程(不是线程)来执行。

标签: multithreading scala apache-spark dataframe concurrency


【解决方案1】:

解决方案

使用this.synchronized{},如下所示

class MyThread extends Runnable{
      def run{
      //some code
         this.synchronized{
            myTable.write.format("parquet").mode("append")
                   .saveAsTable("hdfstable")
         }
      //some code
      }
}
Executors.defaultThreadFactory().newThread(new MyThread).start()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-11-03
    • 2013-11-01
    • 2013-08-08
    • 2017-12-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多