【问题标题】:INSERT HIVE SQL in loop in Spark在 Spark 中的循环中插入 HIVE SQL
【发布时间】:2018-09-27 21:05:29
【问题描述】:

我们正在尝试为 HIVE 运行 INSERT SQL,其中数据来自 Spark 中的数据框。使用的会话有等一切。

有 2 个问题:

问题)即使我们在 forEach 循环中创建了会话,但在尝试使用这两种方法时 INSERT 仍然失败

1) 数据框

2) 直接 Spark SQL

以下是代码(Spark SQL 方法):

import java.time.Instant

import org.apache.spark.sql.{DataFrame, Row, types}
import org.apache.spark.sql.functions.{current_timestamp, first, isnull, lit, max}
import org.apache.spark.sql.types.{StringType, StructField, StructType, TimestampType}

import scala.collection.mutable.ListBuffer

class Controller extends DatabaseServices
  with Loggers {
  val session = createSparkSession(ConfigFactory.load().getString("local.common.spark.app.name"))
  val producer = session.sparkContext.broadcast(KafkaWrapper())

  def doIt(TranIDs: DataFrame): Unit = {
    import session.sqlContext.implicits._

    val TranID = TranIDs
      .withColumnRenamed("TranID", "REFERENCE_TranID")
      .select($"REFERENCE_TranID")
      .union(session.table(BANK_ROLLBACK_TXN_PRODUCER_LOG_VIEW)
        .withColumnRenamed("TranID", "REFERENCE_TranID")
        .select($"REFERENCE_TranID"))
      .where($"REFERENCE_TranID".isNotNull)

    if (TranID.count() == 0) {
      throw new Exception("No rows.")
    }

    val core = session
      .table(BANK_TRANS_MASTER_CORE)
      .withColumnRenamed("TranID", "MASTER_REFERENCE_TranID")
      .withColumnRenamed("CLIENTID", "REF_CLIENT_ID")
      .withColumnRenamed("SUBCLIENTID", "REF_SUBCLIENT_ID")
      .select($"MASTER_REFERENCE_TranID",
        $"TranIDDATE")
      .join(TranID, TranID.col("REFERENCE_TranID") === $"MASTER_REFERENCE_TranID")

    val ref = session
      .table(BANK_RBI_REF_CLIENT)
      .select($"CLIENTID", $"SUBCLIENTID", $"FLAGTRE")
      .join(core, $"CLIENTID" === core.col("REF_CLIENT_ID")
        && $"SUBCLIENTID" === core.col("REF_SUBCLIENT_ID")


    val details = session
      .table(BANK_TRANS_MASTER_DETAILS)
      .select($"TranID",
        $"REALFRAUD",
        $"REALFRAUDDATEBAE",
        $"REALFRAUDYYYYMMDD"
      )
      .join(ref, ref.col("MASTER_REFERENCE_TranID") === $"TranID"
        && $"REALFRAUD" === lit("Y"))
      .where($"TranID".isNotNull
        && $"TranIDDATE".isNotNull)
      .groupBy($"TranID")
      .agg(first($"TranID").as("TranID"),
        first(core("TranIDDATE")).cast("String").as("TranIDDATE"),
        max($"REALFRAUDDATEBAE").as("REALFRAUDDATEBAE"),
        max($"REALFRAUDYYYYMMDD").as("REALFRAUDYYYYMMDD"),
        first($"REALFRAUD").as("REALFRAUD"),
        first($"ABA").as("ABA"))

    details.foreach(row => {


      import scala.collection.JavaConversions._
      val transaction = TxUpdate.newBuilder().setTranID(row.getAs("TranID").toString)
        .setTranIDDATE(row.getAs("TranIDDATE").toString)
        .setAttributes(ListBuffer(
          Attribute.newBuilder.setKey("REALFRAUD").setValue(if (row.getAs("REALFRAUD") != null) row.getAs("REALFRAUD").toString else null).build(),
          Attribute.newBuilder.setKey("REALFRAUDDATEBAE").setValue(if (row.getAs("REALFRAUDDATEBAE") != null) if (row.getAs("REALFRAUDDATEBAE") != null) row.getAs("REALFRAUDDATEBAE").toString else null else null).build(),
          Attribute.newBuilder.setKey("REALFRAUDYYYYMMDD").setValue(if (row.getAs("REALFRAUDYYYYMMDD") != null) row.getAs("REALFRAUDYYYYMMDD").toString else null).build(),
          Attribute.newBuilder.setKey("ABA").setValue(if (row.getAs("ABA") != null) row.getAs("ABA").toString else null).build(),
        .build()

      if (producer.value.sendSync(ConfigFactory.load().getString("local.common.kafka.rollbackKafkaTopicName"),
        transaction.getTranID.toString,
        transaction)) {
        session.sqlContext.sql("insert into " + BANK_ROLLBACK_TXN_PRODUCER_LOG + "(TranID, when_loaded, status) values('" + transaction.getTranID.toString + "', 'current_timestamp()', 'S')")
      } else {
        session.sqlContext.sql("insert into " + BANK_ROLLBACK_TXN_PRODUCER_LOG + "(TranID, when_loaded, status) values('" + transaction.getTranID.toString + "', 'current_timestamp()', 'F')")
      }

    })

  }
}

【问题讨论】:

  • 这是对 Hive 的单例插入。方法不好。
  • 正如我的问题中提到的,我们使用 df.write.insertInto 和 Append 选项也给出了错误。终于让它工作了。感谢您的帮助。

标签: apache-spark hadoop hive apache-spark-sql


【解决方案1】:

这里的错误不清楚。

在较高级别上,您可以使用在 Spark 中启用 hivecontext 的方法,然后使用 append 选项直接将其持久化到 Hive 表中。这将比执行插入操作快得多。流程将是这样的:

步骤 0 - 所有这些都必须在一个 spark 会话中发生。您不需要为每个插入创建多个会话。在某种程度上,这样做是没有意义的。 一种。创建一个数据框,其中包含 Hive 基础表的列。 湾。在火花处理期间,数据帧将其数据最终保存在 Hive 中。 C。使用附加选项启动 Dataframe saveastable

Insert Into Hive

希望这有助于了解您需要如何解决此问题。

【讨论】:

  • 问题不明确怎么回答?
  • 感谢您的回答和帮助。
【解决方案2】:

我们使用了带有 Append 选项的 df.write.insertInto,这会导致错误。终于让它工作了。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-03-09
    • 2019-02-06
    • 1970-01-01
    • 2012-05-20
    • 2016-05-25
    • 1970-01-01
    • 2017-11-29
    • 2015-03-16
    相关资源
    最近更新 更多