Pyspark: come risolvere complicate logiche di dataframe e join

Sep 26 2020

Ho due frame di dati su cui lavorare, il primo assomiglia a questo il seguente df1

df1_schema = StructType([StructField("Date", StringType(), True),\
                              StructField("store_id", StringType(), True),\
                             StructField("warehouse_id", StringType(), True),\
                      StructField("class_id", StringType(), True) ,\
                       StructField("total_time", IntegerType(), True) ])
df_data = [('2020-08-01','110','1','11010',3),('2020-08-02','110','1','11010',2),\
           ('2020-08-03','110','1','11010',3),('2020-08-04','110','1','11010',3),\
            ('2020-08-05','111','1','11010',1),('2020-08-06','111','1','11010',-1)]
rdd = sc.parallelize(df_data)
df1 = sqlContext.createDataFrame(df_data, df1_schema)
df1 = df1.withColumn("Date",to_date("Date", 'yyyy-MM-dd'))
df1.show()

+----------+--------+------------+--------+----------+
|      Date|store_id|warehouse_id|class_id|total_time|
+----------+--------+------------+--------+----------+
|2020-08-01|     110|           1|   11010|         3|
|2020-08-02|     110|           1|   11010|         2|
|2020-08-03|     110|           1|   11010|         3|
|2020-08-04|     110|           1|   11010|         3|
|2020-08-05|     111|           1|   11010|         1|
|2020-08-06|     111|           1|   11010|        -1|
+----------+--------+------------+--------+----------+

Ho calcolato qualcosa chiamato arrival_date

#To calculate the arrival_date
#logic : add the Date + total_time so in first row, 2020-08-01 +3 would give me 2020-08-04 
#if total_time is -1 then return blank
df1= df1.withColumn('arrival_date', F.when(col('total_time') != -1, expr("date_add(date, total_time)"))
        .otherwise(''))
+----------+--------+------------+--------+----------+------------+
|      Date|store_id|warehouse_id|class_id|total_time|arrival_date|
+----------+--------+------------+--------+----------+------------+
|2020-08-01|     110|           1|   11010|         3|  2020-08-04|
|2020-08-02|     110|           1|   11010|         2|  2020-08-04|
|2020-08-03|     110|           1|   11010|         3|  2020-08-06|
|2020-08-04|     110|           1|   11010|         3|  2020-08-07|
|2020-08-05|     111|           1|   11010|         1|  2020-08-06|
|2020-08-06|     111|           1|   11010|        -1|            |
+----------+--------+------------+--------+----------+------------+

e quello che voglio calcolare è questo ..

#to calculate the transit_date
#if arrival_date is same, ex) 2020-08-04 is repeated 2 or more times, then take min("Date") 
#which will be 2020-08-01 otherwise just return the Date ex) 2020-08-07 would just return 2020-08-04
#we need to care about cloth_id too, we have arrival_date = 2020-08-06 repeated 2 times as well but since
#if one of store_id or warehouse_id is different we treat them separately. so at arrival_date = 2020-08-06 at date = 2020-08-03,
##we must return 2020-08-03 
#so we treat them separately when one of (store_id, warehouse_id ) is different. 
#*Note* we dont care about class_id, its not effective.
#if arrival_date = blank then leave it as blank..
#so our df would look something like this.
+----------+--------+------------+--------+----------+------------+------------+
|      Date|store_id|warehouse_id|class_id|total_time|arrival_date|transit_date|
+----------+--------+------------+--------+----------+------------+------------+
|2020-08-01|     110|           1|   11010|         3|  2020-08-04|  2020-08-01|
|2020-08-02|     110|           1|   11010|         2|  2020-08-04|  2020-08-01|
|2020-08-03|     110|           1|   11010|         3|  2020-08-06|  2020-08-03|
|2020-08-04|     110|           1|   11010|         3|  2020-08-07|  2020-08-04|
|2020-08-05|     111|           1|   11010|         1|  2020-08-06|  2020-08-05|
|2020-08-06|     111|           1|   11010|        -1|            |            |
+----------+--------+------------+--------+----------+------------+------------+

Successivamente, ho df2 che assomiglia al seguente ..

#we have another dataframe call it df2

df2_schema = StructType([StructField("Date", StringType(), True),\
                              StructField("store_id", StringType(), True),\
                             StructField("warehouse_id", StringType(), True),\
                             StructField("cloth_id", StringType(), True),\
                      StructField("class_id", StringType(), True) ,\
                       StructField("type", StringType(), True),\
                        StructField("quantity", IntegerType(), True)])
df_data = [('2020-08-01','110','1','M_1','11010','R',5),('2020-08-01','110','1','M_1','11010','R',2),\
           ('2020-08-02','110','1','M_1','11010','C',3),('2020-08-03','110','1','M_1','11010','R',1),\
            ('2020-08-04','110','1','M_1','11010','R',3),('2020-08-05','111','1','M_2','11010','R',5)]
rdd = sc.parallelize(df_data)
df2 = sqlContext.createDataFrame(df_data, df2_schema)
df2 = df2.withColumn("Date",to_date("Date", 'yyyy-MM-dd'))
df2.show()

+----------+--------+------------+--------+--------+----+--------+
|      Date|store_id|warehouse_id|cloth_id|class_id|type|quantity|
+----------+--------+------------+--------+--------+----+--------+
|2020-08-01|     110|           1|     M_1|   11010|   R|       5|
|2020-08-01|     110|           1|     M_1|   11010|   R|       2|
|2020-08-02|     110|           1|     M_1|   11010|   C|       3|
|2020-08-03|     110|           1|     M_1|   11010|   R|       1|
|2020-08-04|     110|           1|     M_1|   11010|   R|       3|
|2020-08-05|     111|           1|     M_2|   11010|   R|       5|
+----------+--------+------------+--------+--------+----+--------+

e ho calcolato quantità2 , questa è solo la somma della quantità dove tipo = R

df2 =df2.groupBy('Date','store_id','warehouse_id','cloth_id','class_id')\
      .agg( F.sum(F.when(col('type')=='R', col('quantity'))\
      .otherwise(col('quantity'))).alias('quantity2')).orderBy('Date')
+----------+--------+------------+--------+--------+---------+
|      Date|store_id|warehouse_id|cloth_id|class_id|quantity2|
+----------+--------+------------+--------+--------+---------+
|2020-08-01|     110|           1|     M_1|   11010|        7|
|2020-08-02|     110|           1|     M_1|   11010|        3|
|2020-08-03|     110|           1|     M_1|   11010|        1|
|2020-08-04|     110|           1|     M_1|   11010|        3|
|2020-08-05|     111|           1|     M_2|   11010|        5|
+----------+--------+------------+--------+--------+---------+

Ora ho df1 e df2. Voglio unirmi in modo che assomigli a questo ... Ho provato qualcosa di simile

df4 = df1.select('store_id','warehouse_id','class_id','arrival_date','transit_date')
df4= df4.filter(" transit_date != '' ")

df4=df4.withColumnRenamed('arrival_date', 'date')

df3 = df2.join(df1, on=['Date','store_id','warehouse_id','class_id'],how='inner').orderBy('Date')
df5 = df3.join(df4, on=['Date','store_id','warehouse_id','class_id'], how='left').orderBy('Date')

ma non credo che questo sia l'approccio corretto .... il risultato df dovrebbe apparire come sotto ..

+----------+--------+------------+--------+--------+---------+----------+------------+------------+
|      Date|store_id|warehouse_id|class_id|cloth_id|quantity2|total_time|arrival_date|transit_date|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+
|2020-08-01|     110|           1|   11010|     M_1|        7|         3|  2020-08-04|        null|
|2020-08-02|     110|           1|   11010|     M_1|        3|         2|  2020-08-04|        null|
|2020-08-03|     110|           1|   11010|     M_1|        1|         3|  2020-08-06|        null|
|2020-08-04|     110|           1|   11010|     M_1|        3|         3|  2020-08-07|  2020-08-01|
|2020-08-05|     111|           1|   11010|     M_2|        5|         1|  2020-08-06|        null|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+

nota che transit_date è andato a dove Date = arrival_dateovviamente il null è sostituito da vuoto.

ULTIMAMENTE, se oggi è 2020-08-04, guarda dove arrival_date == 2020-08-04 e somma la quantità e posizionala a oggi. quindi .... Assomiglierà a questo ... dove store_id = 111, avrà una data separata. non mostrato qui .. quindi la logica deve avere senso anche quando store_id = 111 .. ho appena mostrato l'esempio in cui store_id = 110

Risposte

2 jxc Sep 30 2020 at 01:56

Dalla mia comprensione della tua domanda e dove hai già con quanto segue df1e df2:

df1.orderBy('Date').show()                                           df2.orderBy('Date').show()
+----------+--------+------------+--------+----------+------------+  +----------+--------+------------+--------+--------+---------+
|      Date|store_id|warehouse_id|class_id|total_time|arrival_date|  |      Date|store_id|warehouse_id|cloth_id|class_id|quantity2|
+----------+--------+------------+--------+----------+------------+  +----------+--------+------------+--------+--------+---------+
|2020-08-01|     110|           1|   11010|         3|  2020-08-04|  |2020-08-01|     110|           1|     M_1|   11010|        7|
|2020-08-02|     110|           1|   11010|         2|  2020-08-04|  |2020-08-02|     110|           1|     M_1|   11010|        3|
|2020-08-03|     110|           1|   11010|         3|  2020-08-06|  |2020-08-03|     110|           1|     M_1|   11010|        1|
|2020-08-04|     110|           1|   11010|         3|  2020-08-07|  |2020-08-04|     110|           1|     M_1|   11010|        3|
|2020-08-05|     111|           1|   11010|         1|  2020-08-06|  |2020-08-05|     111|           1|     M_2|   11010|        5|
|2020-08-06|     111|           1|   11010|        -1|            |  +----------+--------+------------+--------+--------+---------+
+----------+--------+------------+--------+----------+------------+

puoi provare i seguenti 5 passaggi:

Passaggio 1: imposta l'elenco dei nomi delle colonne grp_colsper il join:

from pyspark.sql import functions as F
grp_cols = ["Date", "store_id", "warehouse_id", "class_id"]

Step-2: creare DF3 contenente transit_date, che è la data min su ogni combinazione di arrival_date, store_id, warehouse_ide class_id:

df3 = df1.filter('total_time != -1') \
    .groupby("arrival_date", "store_id", "warehouse_id", "class_id") \
    .agg(F.min('Date').alias('transit_date')) \
    .withColumnRenamed("arrival_date", "Date")

df3.orderBy('Date').show()
+----------+--------+------------+--------+------------+
|      Date|store_id|warehouse_id|class_id|transit_date|
+----------+--------+------------+--------+------------+
|2020-08-04|     110|           1|   11010|  2020-08-01|
|2020-08-06|     111|           1|   11010|  2020-08-05|
|2020-08-06|     110|           1|   11010|  2020-08-03|
|2020-08-07|     110|           1|   11010|  2020-08-04|
+----------+--------+------------+--------+------------+

Step-3: configura df4 unendo df2 con df1 e unisciti a df3 a sinistra usando grp_cols, persist df4

df4 = df2.join(df1, grp_cols).join(df3, grp_cols, "left") \
    .withColumn('transit_date', F.when(F.col('total_time') != -1, F.col("transit_date")).otherwise('')) \
    .persist()
_ = df4.count()
df4.orderBy('Date').show()
+----------+--------+------------+--------+--------+---------+----------+------------+------------+
|      Date|store_id|warehouse_id|class_id|cloth_id|quantity2|total_time|arrival_date|transit_date|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+
|2020-08-01|     110|           1|   11010|     M_1|        7|         3|  2020-08-04|        null|
|2020-08-02|     110|           1|   11010|     M_1|        3|         2|  2020-08-04|        null|
|2020-08-03|     110|           1|   11010|     M_1|        1|         3|  2020-08-06|        null|
|2020-08-04|     110|           1|   11010|     M_1|        3|         3|  2020-08-07|  2020-08-01|
|2020-08-05|     111|           1|   11010|     M_2|        5|         1|  2020-08-06|        null|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+

Step-4: calcola sum(quantity2) as wantda df4 per ogni arrival_date+ store_id+ warehouse_id+ class_id+cloth_id

df5 = df4 \
    .groupby("arrival_date", "store_id", "warehouse_id", "class_id", "cloth_id") \
    .agg(F.sum("quantity2").alias("want")) \
    .withColumnRenamed("arrival_date", "Date")
df5.orderBy('Date').show()
+----------+--------+------------+--------+--------+----+
|      Date|store_id|warehouse_id|class_id|cloth_id|want|
+----------+--------+------------+--------+--------+----+
|2020-08-04|     110|           1|   11010|     M_1|  10|
|2020-08-06|     111|           1|   11010|     M_2|   5|
|2020-08-06|     110|           1|   11010|     M_1|   1|
|2020-08-07|     110|           1|   11010|     M_1|   3|
+----------+--------+------------+--------+--------+----+

Passaggio 5: crea il dataframe finale inserendo a sinistra df4 con df5

df_new = df4.join(df5, grp_cols+["cloth_id"], "left").fillna(0, subset=['want'])
df_new.orderBy("Date").show()
+----------+--------+------------+--------+--------+---------+----------+------------+------------+----+
|      Date|store_id|warehouse_id|class_id|cloth_id|quantity2|total_time|arrival_date|transit_date|want|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+----+
|2020-08-01|     110|           1|   11010|     M_1|        7|         3|  2020-08-04|        null|   0|
|2020-08-02|     110|           1|   11010|     M_1|        3|         2|  2020-08-04|        null|   0|
|2020-08-03|     110|           1|   11010|     M_1|        1|         3|  2020-08-06|        null|   0|
|2020-08-04|     110|           1|   11010|     M_1|        3|         3|  2020-08-07|  2020-08-01|  10|
|2020-08-05|     111|           1|   11010|     M_2|        5|         1|  2020-08-06|        null|   0|
+----------+--------+------------+--------+--------+---------+----------+------------+------------+----+
df4.unpersist()
1 Lamanus Sep 27 2020 at 12:19

Ecco per df1,

from pyspark.sql import Window
from pyspark.sql.functions import *
from pyspark.sql.types import *
import builtins as p

df1_schema = StructType(
    [
        StructField('Date',         StringType(),  True),
        StructField('store_id',     StringType(),  True),
        StructField('warehouse_id', StringType(),  True),
        StructField('class_id',     StringType(),  True),
        StructField('total_time',   IntegerType(), True)
    ]
)

df1_data = [
    ('2020-08-01','110','1','11010',3),
    ('2020-08-02','110','1','11010',2),
    ('2020-08-03','110','1','11010',3),
    ('2020-08-04','110','1','11010',3),
    ('2020-08-05','111','1','11010',1),
    ('2020-08-06','111','1','11010',-1)
]


df1 = spark.createDataFrame(df1_data, df1_schema)
df1 = df1.withColumn('Date', to_date('Date'))

df1 = df1.withColumn('arrival_date', when(col('total_time') != -1, expr("date_add(date, total_time)")).otherwise(''))

w = Window.partitionBy('arrival_date', 'store_id', 'warehouse_id').orderBy('Date')
df1 = df1.withColumn('transit_date', when(col('total_time') != -1, first('Date').over(w)).otherwise('')).orderBy('Date')

df1.show()

+----------+--------+------------+--------+----------+------------+------------+
|      Date|store_id|warehouse_id|class_id|total_time|arrival_date|transit_date|
+----------+--------+------------+--------+----------+------------+------------+
|2020-08-01|     110|           1|   11010|         3|  2020-08-04|  2020-08-01|
|2020-08-02|     110|           1|   11010|         2|  2020-08-04|  2020-08-01|
|2020-08-03|     110|           1|   11010|         3|  2020-08-06|  2020-08-03|
|2020-08-04|     110|           1|   11010|         3|  2020-08-07|  2020-08-04|
|2020-08-05|     111|           1|   11010|         1|  2020-08-06|  2020-08-05|
|2020-08-06|     111|           1|   11010|        -1|            |            |
+----------+--------+------------+--------+----------+------------+------------+

e df2 come hai fatto tu,

df2_schema = StructType(
    [
        StructField('Date',         StringType(),  True),
        StructField('store_id',     StringType(),  True),
        StructField('warehouse_id', StringType(),  True),
        StructField('cloth_id',     StringType(),  True),
        StructField('class_id',     StringType(),  True),
        StructField('type',         StringType(),  True),
        StructField('quantity',     IntegerType(), True)
    ]
)

df2_data = [
    ('2020-08-01','110','1','M_1','11010','R',5),
    ('2020-08-01','110','1','M_1','11010','R',2),
    ('2020-08-02','110','1','M_1','11010','C',3),
    ('2020-08-03','110','1','M_1','11010','R',1),
    ('2020-08-04','110','1','M_1','11010','R',3),
    ('2020-08-05','111','1','M_2','11010','R',5)
]

df2 = spark.createDataFrame(df2_data, df2_schema)
df2 = df2.withColumn('Date', to_date('Date'))

df2 = df2.groupBy('Date', 'store_id', 'warehouse_id', 'cloth_id', 'class_id') \
        .agg(
            sum(
                when(col('type') == 'R', col('quantity')).otherwise(0)
            ).alias('quantity2')
        ).orderBy('Date')

df2.show()

+----------+--------+------------+--------+--------+---------+
|      Date|store_id|warehouse_id|cloth_id|class_id|quantity2|
+----------+--------+------------+--------+--------+---------+
|2020-08-01|     110|           1|     M_1|   11010|        7|
|2020-08-02|     110|           1|     M_1|   11010|        0|
|2020-08-03|     110|           1|     M_1|   11010|        1|
|2020-08-04|     110|           1|     M_1|   11010|        3|
|2020-08-05|     111|           1|     M_2|   11010|        5|
+----------+--------+------------+--------+--------+---------+

e infine il risultato del join.

df3 = df1.filter('total_time != -1') \
  .join(df2, on=['Date', 'store_id', 'warehouse_id', 'class_id'], how='left') \
  .drop('Date', 'total_time', 'cloth_id') \
  .withColumnRenamed('arrival_date', 'Date')

df4 = df1.drop('transit_date') \
  .join(df3, on=['Date', 'store_id', 'warehouse_id', 'class_id'], how='left') \
  .groupBy('Date', 'store_id', 'warehouse_id', 'class_id', 'arrival_date', 'transit_date') \
  .agg(sum('quantity2').alias('want')) \
  .orderBy('Date')

df4.show()

+----------+--------+------------+--------+------------+------------+----+
|      Date|store_id|warehouse_id|class_id|arrival_date|transit_date|want|
+----------+--------+------------+--------+------------+------------+----+
|2020-08-01|     110|           1|   11010|  2020-08-04|        null|null|
|2020-08-02|     110|           1|   11010|  2020-08-04|        null|null|
|2020-08-03|     110|           1|   11010|  2020-08-06|        null|null|
|2020-08-04|     110|           1|   11010|  2020-08-07|  2020-08-01|   7|
|2020-08-05|     111|           1|   11010|  2020-08-06|        null|null|
|2020-08-06|     111|           1|   11010|            |  2020-08-05|   5|
+----------+--------+------------+--------+------------+------------+----+