【问题标题】:How to update few records in Spark如何在 Spark 中更新几条记录
【发布时间】:2018-07-15 15:03:19
【问题描述】:

我在 Scala 中有以下用于 spark 的程序:

val dfA = sqlContext.sql("select * from employees where id in ('Emp1', 'Emp2')" )
val dfB = sqlContext.sql("select * from employees where id not in ('Emp1', 'Emp2')" )
val dfN = dfA.withColumn("department", lit("Finance"))
val dfFinal = dfN.unionAll(dfB)
dfFinal.registerTempTable("intermediate_result")

dfA.unpersist
dfB.unpersist
dfN.unpersist
dfFinal.unpersist

val dfTmp = sqlContext.sql("select * from intermediate_result")
dfTmp.write.mode("overwrite").format("parquet").saveAsTable("employees")
dfTmp.unpersist

当我尝试保存它时,我收到以下错误:

org.apache.spark.sql.AnalysisException: 无法覆盖正在读取的表employees。; 在 org.apache.spark.sql.execution.datasources.PreWriteCheck.failAnalysis(rules.scala:106) 在 org.apache.spark.sql.execution.datasources.PreWriteCheck$$anonfun$apply$3.apply(rules.scala:182) 在 org.apache.spark.sql.execution.datasources.PreWriteCheck$$anonfun$apply$3.apply(rules.scala:109) 在 org.apache.spark.sql.catalyst.trees.TreeNode.foreach(TreeNode.scala:111) 在 org.apache.spark.sql.execution.datasources.PreWriteCheck.apply(rules.scala:109) 在 org.apache.spark.sql.execution.datasources.PreWriteCheck.apply(rules.scala:105) 在 org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$2.apply(CheckAnalysis.scala:218) 在 org.apache.spark.sql.catalyst.analysis.CheckAnalysis$$anonfun$checkAnalysis$2.apply(CheckAnalysis.scala:218) 在 scala.collection.immutable.List.foreach(List.scala:318)

我的问题是:

  1. 我更改两名员工部门的方法是否正确
  2. 为什么我在释放 DataFrames 后会收到此错误

【问题讨论】:

    标签: scala apache-spark dataframe hive


    【解决方案1】:

    我改变两名员工的部门的方法是否正确

    事实并非如此。只是重复在 Stack Overflow 上多次说过的话 - Apache Spark 不是数据库。它不是为细粒度更新而设计的。如果您的项目需要这样的操作,请使用 Hadoop 上的众多数据库之一。

    为什么我在发布 DataFrame 后会收到此错误

    因为你没有。您所做的只是为执行计划添加一个名称。检查点是最接近“释放”的东西,但你真的不希望在破坏性操作过程中失去执行者的情况。

    您可以写入临时目录、删除输入并移动临时文件,但实际上 - 只需使用适合该工作的工具。

    【讨论】:

    • 我该如何解决?
    【解决方案2】:

    以下是您可以尝试的方法。

    您可以使用 saveAsTable api 将其写入另一个表,而不是使用 registertemptable api

    dfFinal.write.mode("overwrite").saveAsTable("intermediate_result")
    

    然后,写入员工表

     val dy = sqlContext.table("intermediate_result")
      dy.write.mode("overwrite").insertInto("employees")
    

    最后,删除中间结果表。

    【讨论】:

      【解决方案3】:

      我会这样处理,

      >>> df = sqlContext.sql("select * from t")
      >>> df.show()
      +-------------+---------------+
      |department_id|department_name|
      +-------------+---------------+
      |            2|        Fitness|
      |            3|       Footwear|
      |            4|        Apparel|
      |            5|           Golf|
      |            6|       Outdoors|
      |            7|       Fan Shop|
      +-------------+---------------+
      

      为了模仿您的流程,我创建了 2 个数据帧,执行 union 并写回 同一张表 t(在本例中故意删除 department_id = 4

      >>> df1 = sqlContext.sql("select * from t where department_id < 4")
      >>> df2 = sqlContext.sql("select * from t where department_id > 4")
      >>> df3 = df1.unionAll(df2)
      >>> df3.registerTempTable("df3")
      >>> sqlContext.sql("insert overwrite table t select * from df3")
      DataFrame[]  
      >>> sqlContext.sql("select * from t").show()
      +-------------+---------------+
      |department_id|department_name|
      +-------------+---------------+
      |            2|        Fitness|
      |            3|       Footwear|
      |            5|           Golf|
      |            6|       Outdoors|
      |            7|       Fan Shop|
      +-------------+---------------+
      

      【讨论】:

      • 第 6 行的“DataFrame[]”是做什么的?还是打字错误?
      • 没什么。 sqlContext.sql 只是返回了一个 dataframe,在这种情况下没有分配给任何变量。
      【解决方案4】:

      假设它是您正在读取和覆盖的配置单元表

      请按如下方式将时间戳引入hive表位置

          create table table_name (
        id                int,
        dtDontQuery       string,
        name              string
      )
       Location hdfs://user/table_name/timestamp
      

      由于无法覆盖,我们会将输出文件写入新位置。

      使用数据框 Api 将数据写入新位置

      df.write.orc(hdfs://user/xx/tablename/newtimestamp/)
      

      写入数据后,将配置单元表位置更改为新位置

      Alter table tablename set Location hdfs://user/xx/tablename/newtimestamp/
      

      【讨论】:

      • 我真的看不懂
      • 是的,一切都在 hive 上。而且该表有很多列。我从哪里发出这些命令?直线?还是来自火花工作?
      猜你喜欢
      • 1970-01-01
      • 2018-11-21
      • 2014-06-21
      • 1970-01-01
      • 2016-08-13
      • 1970-01-01
      • 2019-02-18
      • 2012-09-29
      • 1970-01-01
      相关资源
      最近更新 更多