【问题标题】:How to convert complex SQL query to spark-dataframe using python or Scala如何使用 python 或 Scala 将复杂的 SQL 查询转换为 spark-dataframe
【发布时间】:2021-02-01 07:44:30
【问题描述】:

我在 spark 中使用 sqlcontext 完成了一次转换,但我只想使用 Spark Data frame 编写相同的查询。该查询包括连接操作和 SQL 的 case 语句。 sql查询编写如下:

refereshLandingData=spark.sql( "select a.Sale_ID, a.Product_ID,"
                           "CASE "
                           "WHEN (a.Quantity_Sold IS NULL) THEN b.Quantity_Sold "
                           "ELSE a.Quantity_Sold "
                           "END AS Quantity_Sold, "
                           "CASE "
                           "WHEN (a.Vendor_ID IS NULL) THEN b.Vendor_ID "
                           "ELSE a.Vendor_ID "
                           "END AS Vendor_ID, "
                           "a.Sale_Date, a.Sale_Amount, a.Sale_Currency "
                           "from landingData a left outer join preHoldData b on a.Sale_ID = b.Sale_ID" )

现在我想要 scala 和 python 中的 spark 数据帧中的等效代码。我尝试了一些代码,但它的
不工作。我试过的代码如下:

joinDf=landingData.join(preHoldData,landingData['Sale_ID']==preHoldData['Sale_ID'],'left_outer')

joinDf.withColumn\
('QuantitySold',pf.when(pf.col(landingData('Quantity_Sold')).isNull(),pf.col(preHoldData('Quantity_Sold')))
.otherwise(pf.when(pf.col(preHoldData('Quantity_Sold')).isNull())),
 pf.col(landingData('Quantity_Sold'))).show()

在上面的代码中,连接完成得很完美,但案例条件不起作用。 我得到--> TypeError: 'DataFrame' object is not callable 我正在使用 spark 2.3.2 版本和 python 3.7 以及类似的 scala 2.11,以防 spark-scala 请任何人建议我任何等效的代码或指南!

【问题讨论】:

  • 检查您的 Python 代码,因为您正在尝试从数据框实例调用函数

标签: python scala apache-spark pyspark apache-spark-sql


【解决方案1】:

这是一个 scala 解决方案: 假设 landingDatapreHoldData 是您的数据框


 val landingDataDf = landingData.withColumnRenamed("Quantity_Sold","Quantity_Sold_ld")
 val preHoldDataDf = preHoldData.withColumnRenamed("Quantity_Sold","Quantity_Sold_phd")

 val joinDf = landingDataDf.join(preHoldDataDf, Seq("Sale_ID"))


 joinDf
 .withColumn("Quantity_Sold",
    when(col("Quantity_Sold_ld").isNull , col("Quantity_Sold_phd")).otherwise(col("Quantity_Sold_ld"))
 ). drop("Quantity_Sold_ld","Quantity_Sold_phd")

您可以对 Vendor_id 执行相同的操作

您的代码的问题是,您无法在 withColumn 操作中引用其他/旧数据框名称。它必须来自您正在操作的数据框。

【讨论】:

  • 此代码工作正常,但问题是数据框具有相同的列名,我们必须重命名该列中的每一列,否则无法选择所需的列。我们只希望在这种情况下选择列而不是全部.
  • @AliBinmazi 您可以轻松删除列。编辑答案以删除不需要的列。
  • 现在,当我们添加 drop() 时,这段代码可以正常工作。谢谢@Sanket9394
  • @AliBinmazi 很酷。如果您发现它有用,请考虑接受它作为答案:)
【解决方案2】:

下面的代码可以在 scala 和 python 上运行,你可以稍微调整一下。

val preHoldData = spark.table("preHoldData").alias("a")
val landingData = spark.table("landingData").alias("b")

landingData.join(preHoldData,Seq("Sale_ID"),"leftouter")
.withColumn("Quantity_Sold",when(col("a.Quantity_Sold").isNull, col("b.Quantity_Sold")).otherwise(col("a.Quantity_Sold")))
.withColumn("Vendor_ID",when(col("a.Vendor_ID").isNull, col("b.Vendor_ID")).otherwise(col("a.Vendor_ID")))
.select(col("a.Sale_ID"),col("a.Product_ID"),col("Quantity_Sold"),col("Vendor_ID"),col("a.Sale_Date"),col("a.Sale_Amount"),col("a.Sale_Currency"))

【讨论】:

  • 它在使用column 函数之前工作正常,但是当我们应用.select 函数时,col("Quantity_Sold") 会出现歧义问题。我尝试重命名它但仍然无法正常工作。
  • 好的,检查@Sanket9394他的解决方案..他已经解释清楚了...... :)
猜你喜欢
  • 2015-04-04
  • 2020-12-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-29
相关资源
最近更新 更多