【问题标题】:Update with inner join using spark dataframe/dataset/RDD使用 spark dataframe/dataset/RDD 更新内部连接
【发布时间】:2018-08-09 10:03:29
【问题描述】:

我正在将 ms sql server 查询的逻辑转换为 spark。要转换的查询如下:

Update enc set PrUid=m.OriginalPrUid
FROM CachePatDemo enc 
inner join #MergePreMap m on enc.PrUid=m.NewPrUid
WHERE StatusId is null

我正在使用数据框进行转换,并且我的两个数据框中有两个表,我将它们作为内部连接加入。我需要找到一种方法来获取表 1 的所有列和更新的列(这在两个表中都很常见)。

我试过用这个:

val result = CachePatDemo.as("df123").
  join(MergePreMap.as("df321"), CachePatDemo("prUid") === MergePreMap("prUid"),"inner").where("StatusId is null")
  select($"df123.pId", 
         $"df321.provFname".as("firstName"), 
         $"df123.lastName", 
         $"df123.prUid")

这似乎没有解决我的问题。有人可以帮忙吗?

【问题讨论】:

  • 两个输入数据框的架构和输出数据框的架构应该可以帮助您快速获得答案
  • 代码中的newdfdf2 指的是什么?澄清(请不要在 cmets 中,编辑您的帖子)
  • 您介意请检查您的问题吗?我相信您提供的代码存在一些不一致之处。然而,我也没有看到您在 spark-sql 查询中添加了 null 条件的位置。
  • @AlexSavitsky 更新了问题

标签: sql sql-server apache-spark apache-spark-sql


【解决方案1】:

在 Spark 2.1 上有效

case class TestModel(x1: Int, x2: String, x3: Int)

object JoinDataFrames extends App {
  import org.apache.spark.sql.{DataFrame, SparkSession}
  val spark = SparkSession.builder.appName("GroupOperations").master("local[2]").enableHiveSupport.getOrCreate

  import spark.implicits._
  import org.apache.spark.sql.functions._

  val list1 = (3 to 10).toList.map(i => new TestModel(i, "This is df1 " + i, i * 3))
  val list2 = (0 to 5).toList.map(i => new TestModel(i, "This is df2 " + i, i * 13))
  val df1: DataFrame = spark.sqlContext.createDataFrame[TestModel](list1)
  val df2: DataFrame = spark.sqlContext.createDataFrame[TestModel](list2)
  val res = df1.join(df2, Seq("x1"), "inner")
  println("from DF1")
  res.select(df1("x2")).show()
  println("from DF2")
  res.select(df2("x2")).show()
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-09-28
    • 2021-11-19
    • 2019-06-16
    • 2020-08-26
    • 2021-11-05
    • 1970-01-01
    • 2017-04-01
    • 2017-02-03
    相关资源
    最近更新 更多