【问题标题】:How to reduce multiple joins in spark如何减少火花中的多个连接
【发布时间】:2021-05-16 19:59:08
【问题描述】:

我使用 Spark 2.4.1 来计算我的数据框的一些比率。

我需要通过加入元数据帧(即resDs)来找到给定数据帧(df_data)中不同比率的比率因子。

我通过使用具有不同连接条件的三个不同连接(即joinedDs、joinedDs2、joinedDs3)来获得这些比率因子(即ratio_1_factor、ratio_2_factor 和ratio_3_factor)

还有其他方法可以减少连接数吗?让它发挥最佳效果?

您可以在以下公共 URL 中找到完整的示例数据。

https://databricks-prod-cloudfront.cloud.databricks.com/public/4027ec902e239c93eaaa8714f173bcfc/1165111237342523/3521103084252405/7035720262824085/latest.html

如何在when子句中处理多步而不是单步:

.withColumn("step_1_ratio_1", (col("ratio_1").minus(lit(0.00000123))).cast(DataTypes.DoubleType)) // step-2
      .withColumn("step_2_ratio_1", (col("step_1_ratio_1").multiply(lit(0.02))).cast(DataTypes.DoubleType)) //step-3
      .withColumn("step_3_ratio_1", (col("step_2_ratio_1").divide(col("step_1_ratio_1"))).cast(DataTypes.DoubleType)) //step-4
      .withColumn("ratio_1_factor", (col("ratio_1_factor")).cast(DataTypes.DoubleType)) //step-5

即“ratio_1_factor”根据数据框中的其他列计算,df_data

这些步骤 -2,3,4 也用于其他 ratio_factors 计算。即 ratio_2_factor, ratio_2_factor 这应该怎么处理?

【问题讨论】:

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


    【解决方案1】:

    您可以加入一次并使用聚合中的max 和when 函数计算ratio_1_factor、ratio_2_factor 和ratio_3_factor 列:

    val joinedDs = df_data.as("aa")
      .join(
        broadcast(resDs.as("bb")),
        col("aa.g_date").between(col("bb.start_date"), col("bb.end_date"))
      )
      .groupBy("item_id", "g_date", "ratio_1", "ratio_2", "ratio_3")
      .agg(
        max(when(
            col("aa.ratio_1").between(col("bb.A"), col("bb.A_lead")),
            col("ratio_1").multiply(lit(0.1))
          )
        ).cast(DoubleType).as("ratio_1_factor"),
        max(when(
            col("aa.ratio_2").between(col("bb.A"), col("bb.A_lead")),
            col("ratio_2").multiply(lit(0.2))
          )
        ).cast(DoubleType).as("ratio_2_factor"),
        max(when(
            col("aa.ratio_3").between(col("bb.A"), col("bb.A_lead")),
            col("ratio_3").multiply(lit(0.3))
          )
        ).cast(DoubleType).as("ratio_3_factor")
      )
    
    
    joinedDs.show(false)
    
    //+-------+----------+---------+-----------+-----------+---------------------+--------------+--------------+
    //|item_id|g_date    |ratio_1  |ratio_2    |ratio_3    |ratio_1_factor       |ratio_2_factor|ratio_3_factor|
    //+-------+----------+---------+-----------+-----------+---------------------+--------------+--------------+
    //|50312  |2016-01-04|0.0456646|0.046899415|0.046000415|0.0045664600000000005|0.009379883   |0.0138001245  |
    //+-------+----------+---------+-----------+-----------+---------------------+--------------+--------------+
    

    【讨论】:

    • 非常感谢,我还有一个疑问,在“when”子句中,如果我需要处理多步而不是一个,即 col("ratio_x").multiply(lit(0. x)) ,如何处理。这些多步骤,即新列值也用于其他比率列,即 ratio_1_factor 、 ratio_2_factor 和 ratio_3_factor 。我已经更新了新的查询,请您指教。
    • @BdEngineer 如果你这样写有什么问题:.withColumn("ratio_1_factor", (col("ratio_1") - lit(0.00000123)) * lit(0.02) / (col("ratio_1") - lit(0.00000123)))?像这样与when 一起使用没有问题。你只需要简化表达式,when表达式背后的逻辑是一样的。
    • 这样写查询的步骤/方法是什么,应该如何培养像你这样的技能?你是怎么写的……我仍然在努力正确使用 agg 函数……我缺少一些基本的东西……请帮助……先生,我有一个像这样的用例stackoverflow.com/questions/66635632/…可以你建议如何处理它?
    猜你喜欢
    • 2017-05-13
    • 1970-01-01
    • 1970-01-01
    • 2016-06-05
    • 2019-06-10
    • 1970-01-01
    • 2018-01-08
    • 2022-01-15
    • 1970-01-01
    相关资源
    最近更新 更多