【问题标题】:Spark - How to update value using data frame ScalaSpark - 如何使用数据框 Scala 更新值
【发布时间】:2016-08-31 00:39:34
【问题描述】:

我有两个结构如下的文件

文件 1

gnk_id, matchId, timestamp

文件 2

gnk_matchid, matchid

如果file1.gnk_id = file2.gnk_machid,我想用文件2中matchid的值更新文件1中gnk_id的值。

为此,我在 Spark 中创建了两个数据框。我想知道我们是否可以更新 Spark 中的值?如果没有,是否有任何解决方法可以提供更新的最终文件?

更新

我做了这样的事情

case class GnkMatchId(gnk: String, gnk_matchid: String)
case class MatchGroup(gnkid: String, matchid: String, ts: String)

val gnkmatchidRDD = sc.textFile("000000000001").map(_.split(',')).map(x => (x(0),x(1)) )

val gnkmatchidDF = gnkmatchidRDD.map( x => GnkMatchId(x._1,x._2) ).toDF()

val matchGroupMr = sc.textFile("part-00000").map(_.split(',')).map(x => (x(0),x(1),x(2)) ).map( f => MatchGroup(f._1,f._2,f._3.toString) ).toDF()

val matchgrp_joinDF = matchGroupMr.join(gnkmatchidDF,matchGroupMr("gnkid") === gnkmatchidDF("gnk_matchid"),"left_outer")

matchgrp_joinDF.map(x => if(x.getAs[String]("gnk_matchid").length != 0 ) {MatchGroup(x.getAs[String]("gnk_matchid"), x.getAs[String]("matchid"),x.getAs[String]("ts"))} else {MatchGroup(x.getAs[String]("gnkid"), x.getAs[String]("matchid"),x.getAs[String]("ts"))}).toDF().show()

但在最后一步,NULLpointerEXception 失败了

【问题讨论】:

标签: apache-spark apache-spark-sql


【解决方案1】:

DataFrame 是基于RDD 的,所以你不能更新其中的值。

但您可以通过添加新列来更新值 withColumn

在你的情况下,你可以通过 joinwithColumn 使用 UDF:

// df1: your File1
// +------+-------+---+
// |gnk_id|matchId| ts|
// +------+-------+---+
// |     1|     10|100|
// |     2|     20|200|
// +------+-------+---+

// df2: your File2
// +-----------+-------+
// |gnk_matchid|matchid|
// +-----------+-------+
// |          1|   1000|
// |          3|   3000|
// +-----------+-------+

// UDF: choose values from matchid or gnk_id for the new column
val myUDF = udf[Integer,Integer,Integer]((df2_matchid: Integer, df1_gnk_id: Integer) => {
  if (df2_matchid == null) df1_gnk_id
  else df2_matchid
})

df1.join(df2, $"gnk_id"===$"gnk_matchid", "left_outer")
  .select($"df1.*", $"df2.matchid" as "matchid2")
  .withColumn("gnk_id", myUDF($"matchid2", $"gnk_id"))
  .drop($"matchid2")
  .show()

输出如下:

+------+-------+---+
|gnk_id|matchId| ts|
+------+-------+---+
|  1000|     10|100|
|     2|     20|200|
+------+-------+---+

【讨论】:

    【解决方案2】:

    这可能是您正在寻找的加入。假设您的数据框位于 file1file2 中,您可以尝试以下操作:

    val result = file1
      .join(file2, file1("matchId") === file2("matchid"))
      .select(
        col("gnk_matchid").as("gnk_id"),
        col("matchId"),
        col("timestamp")
      )
    

    【讨论】:

      【解决方案3】:

      这取决于您使用的数据源是否支持。

      • 与 Hive 你can
      • 对于文件,您只需过滤掉行,进行更改并重新添加即可。

      希望对您有所帮助。

      2018/04/11 更新链接

      【讨论】:

      • 您的第一个链接已失效,而且我质疑您是否可以在 Hive 表上使用 Spark SQL 进行更新,因为有一张支持更新 Hive(事务)表的开放票:issues.apache.org/jira/browse/SPARK-15348
      • 来自这本书。 Hive 解决方案只是连接它不会更改或更改记录的文件。可以使用 ORC 格式更新 Hive 中的数据使用 Hive 中的事务表以及插入、更新、删除,它会定期自动为您“连接”。目前这仅适用于 orc.format 的表(存储为 orc)或者,使用 Hbase 和 Phoenix 作为顶部的 SQL 层 Hive 最初不是为更新而设计的,因为它是.purely 仓库集中的,最新的可以更新, 以事务方式删除等。
      • 确实可以使用 Hive 对 ORC 格式的 Hive 事务表进行更新。但是,Spark SQL 不支持这一点:目前无法在 Hive 事务表上使用 Spark SQL 进行更新,如我提到的票证中所述issues.apache.org/jira/browse/SPARK-15348
      【解决方案4】:

      实现的最简单方法,下面的代码读取每个批次的维度数据文件夹,但请记住新的维度数据值(在我的例子中为国家名称)必须是一个新文件。

      以下流+批量连接的解决方案

      package com.databroccoli.streaming.dimensionupateinstreaming
      
      import org.apache.log4j.{Level, Logger}
      import org.apache.spark.sql.{DataFrame, ForeachWriter, Row, SparkSession}
      import org.apache.spark.sql.functions.{broadcast, expr}
      import org.apache.spark.sql.types.{StringType, StructField, StructType, TimestampType}
      
      object RefreshDimensionInStreaming {
      
        def main(args: Array[String]) = {
      
          @transient lazy val logger: Logger = Logger.getLogger(getClass.getName)
      
          Logger.getLogger("akka").setLevel(Level.WARN)
          Logger.getLogger("org").setLevel(Level.ERROR)
          Logger.getLogger("com.amazonaws").setLevel(Level.ERROR)
          Logger.getLogger("com.amazon.ws").setLevel(Level.ERROR)
          Logger.getLogger("io.netty").setLevel(Level.ERROR)
      
          val spark = SparkSession
            .builder()
            .master("local")
            .getOrCreate()
      
          val schemaUntyped1 = StructType(
            Array(
              StructField("id", StringType),
              StructField("customrid", StringType),
              StructField("customername", StringType),
              StructField("countrycode", StringType),
              StructField("timestamp_column_fin_1", TimestampType)
            ))
      
          val schemaUntyped2 = StructType(
            Array(
              StructField("id", StringType),
              StructField("countrycode", StringType),
              StructField("countryname", StringType),
              StructField("timestamp_column_fin_2", TimestampType)
            ))
      
          val factDf1 = spark.readStream
            .schema(schemaUntyped1)
            .option("header", "true")
            .csv("src/main/resources/broadcasttest/fact")
      
          var countryDf: Option[DataFrame] = None: Option[DataFrame]
      
          def updateDimensionDf() = {
            val dimDf2 = spark.read
              .schema(schemaUntyped2)
              .option("header", "true")
              .csv("src/main/resources/broadcasttest/dimension")
      
            if (countryDf != None) {
              countryDf.get.unpersist()
            }
      
            countryDf = Some(
              dimDf2
                .withColumnRenamed("id", "id_2")
                .withColumnRenamed("countrycode", "countrycode_2"))
      
            countryDf.get.show()
          }
      
          factDf1.writeStream
            .outputMode("append")
            .foreachBatch { (batchDF: DataFrame, batchId: Long) =>
              batchDF.show(10)
      
              updateDimensionDf()
      
              batchDF
                .join(
                  countryDf.get,
                  expr(
                    """
            countrycode_2 = countrycode 
            """
                  ),
                  "leftOuter"
                )
                .show
      
            }
            .start()
            .awaitTermination()
      
        }
      
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2021-11-23
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-10-09
        • 2016-04-04
        相关资源
        最近更新 更多