【问题标题】:Daily forecast on a PySpark dataframePySpark 数据框的每日预测
【发布时间】:2021-05-02 13:55:17
【问题描述】:

我在 PySpark 中有以下数据框:

DT_BORD_REF:月份的日期列
REF_DATE:分隔过去和未来日期的日期参考
PROD_ID:产品 ID
COMPANY_CODE:公司 ID
CUSTOMER_CODE:客户 ID
MTD_WD:本月至今的工作日计数(日期 = DT_BORD_REF)
QUANTITY:售出的商品数量
QTE_MTD:本月至今的商品数量

+-------------------+-------------------+-----------------+------------+-------------+-------------+------+--------+-------+
|        DT_BORD_REF|           REF_DATE|          PROD_ID|COMPANY_CODE|CUSTOMER_CODE|COUNTRY_ALPHA|MTD_WD|QUANTITY|QTE_MTD|
+-------------------+-------------------+-----------------+------------+-------------+-------------+------+--------+-------+
|2020-11-02 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     1|     4.0|    4.0|
|2020-11-05 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     3|    null|    4.0|
|2020-11-06 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     4|    null|    4.0|
|2020-11-09 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     5|    null|    4.0|
|2020-11-10 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     6|    null|    4.0|
|2020-11-11 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     7|    null|    4.0|
|2020-11-12 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     8|    null|    4.0|
|2020-11-13 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|     9|    null|    4.0|
|2020-11-16 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    10|    null|    4.0|
|2020-11-17 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    11|    null|    4.0|
|2020-11-18 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    12|    null|    4.0|
|2020-11-19 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    13|    null|    4.0|
|2020-11-20 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    14|    null|    4.0|
|2020-11-23 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    15|    null|    4.0|
|2020-11-24 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    16|    null|    4.0|
|2020-11-25 00:00:00|2020-11-04 00:00:00|          0000043|         503|     KDAI3982|          RUS|    17|    null|    4.0|

对于DT_BORD_REF < REF_DATE,所有行都是实际销售额,不一定每个工作日都发生。有时也发生在非工作日。

DT_BORD_REF >= REF_DATE 没有销售(这是未来)

目标是使用以下公式预测所有未来行的销售额:QTE_MTD/MTD_WD 根据每个产品、客户和国家/地区的REF_DATE 计算得出。

使用窗口函数从 QUANTITY 列计算 QTE_MTD。我需要将 MTD_WD 划分为 REF_DATE,在这个例子中是 3

如何在REF_DATE 上添加带有MTD_WD 的列,按产品、客户和国家/地区进行分区?

换句话说,当满足每个产品、客户和国家/地区的条件 DT_BORD_REF > REF_DATE(同样,在此示例中为 3)时,我需要添加一个第一次出现 MTD_WD 的列。

此数据集包含数百万行,用于不同的产品、客户和国家/地区 工作日按国家/地区提供

希望清楚:)

【问题讨论】:

    标签: sql apache-spark pyspark apache-spark-sql window-functions


    【解决方案1】:

    您可以将firstignorenulls=True 一起使用,并将when 与适当的条件一起使用,以获取第一个MTD_WD 其中DT_BORD_REF > REF_DATE

    from pyspark.sql import functions as F, Window
    
    df2 = df.withColumn(
        'val',
        F.first(
            F.when(
                F.col('DT_BORD_REF') > F.col('REF_DATE'),
                F.col('MTD_WD')
            ), 
            ignorenulls=True
        ).over(
            Window.partitionBy('PROD_ID','COMPANY_CODE','CUSTOMER_CODE','COUNTRY_ALPHA')
                  .orderBy('DT_BORD_REF')
                  .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
        )
    )
    

    【讨论】:

      猜你喜欢
      • 2021-11-30
      • 2017-07-31
      • 2020-07-12
      • 2017-05-20
      • 2021-09-20
      • 1970-01-01
      • 2020-02-20
      • 2016-04-21
      • 1970-01-01
      相关资源
      最近更新 更多