Row number with lag condition on multiple columns

Viewed 422

I would like to create row number that is partitioned by ACCOUNT, NAME and TYPE.

I tried dense rank and row number. However, I need all initial records that contain changes in any of those columns

    df = spark.createDataFrame(
       [
    ('20190910', 'A1', 'Linda', 'b2c'),
    ('20190911', 'A1', 'Tom', 'consultant'),
    ('20190912', 'A1', 'John', 'b2c'),
    ('20190913', 'A1', 'Tom', 'consultant'),
    ('20190914', 'A1', 'Tom', 'consultant'),
    ('20190915', 'A1', 'Linda', 'consultant'),
    ('20190916', 'A1', 'Linda', 'b2c'),
    ('20190917', 'B1', 'John', 'b2c'),
    ('20190916', 'B1', 'John', 'consultant'),
    ('20190910', 'B1', 'Linda', 'b2c'),
    ('20190911', 'B1', 'John', 'b2c'),
    ('20190915', 'C1', 'John', 'consultant'),
    ('20190916', 'C1', 'Linda', 'consultant'),
    ('20190917', 'C1', 'John', 'b2c'),
    ('20190916', 'C1', 'RJohn', 'consultant'),
    ('20190910', 'C1', 'Tom', 'b2c'),
    ('20190911', 'C1', 'John', 'b2c'),
     ],
    ['Event_date', 'account', 'name', 'type']
     )

Expected outcome:

Event_date account name type row_number
20190910 A1 Linda b2c 1
20190911 A1 Tom consultant 1
20190912 A1 John b2c 1
20190913 A1 Tom consultant 2
20190914 A1 Tom consultant 3
20190915 A1 Linda consultant 1
20190916 A1 Linda b2c 2
20190917 B1 John b2c 1
20190916 B1 John consultant 1
20190910 B1 Linda b2c 2
20190911 B1 John b2c 3
20190915 C1 John consultant 1
20190916 C1 Linda consultant 1
20190917 C1 John b2c 1
20190916 C1 John consultant 2
20190910 C1 Tom b2c 1
20190911 C1 John b2c 2
2 Answers

You could create a Window and partition it by account, name, type and then row_number over it.

Example:

spark = SparkSession.builder.getOrCreate()
df = spark.createDataFrame(
    [
        ("20190910", "A1", "Linda", "b2c"),
        ("20190911", "A1", "Tom", "consultant"),
        ("20190912", "A1", "John", "b2c"),
        ("20190913", "A1", "Tom", "consultant"),
        ("20190914", "A1", "Tom", "consultant"),
        ("20190915", "A1", "Linda", "consultant"),
        ("20190916", "A1", "Linda", "b2c"),
        ("20190917", "B1", "John", "b2c"),
        ("20190916", "B1", "John", "consultant"),
        ("20190910", "B1", "Linda", "b2c"),
        ("20190911", "B1", "John", "b2c"),
        ("20190915", "C1", "John", "consultant"),
        ("20190916", "C1", "Linda", "consultant"),
        ("20190917", "C1", "John", "b2c"),
        ("20190916", "C1", "RJohn", "consultant"),
        ("20190910", "C1", "Tom", "b2c"),
        ("20190911", "C1", "John", "b2c"),
    ],
    ["Event_date", "account", "name", "type"],
)
w = Window.partitionBy("account", "name", "type").orderBy("Event_date")
df = df.withColumn("row_number", F.row_number().over(w)).orderBy("Event_date")

Result:

+----------+-------+-----+----------+----------+                                
|Event_date|account|name |type      |row_number|
+----------+-------+-----+----------+----------+
|20190912  |A1     |John |b2c       |1         |
|20190911  |A1     |Tom  |consultant|1         |
|20190913  |A1     |Tom  |consultant|2         |
|20190914  |A1     |Tom  |consultant|3         |
|20190915  |A1     |Linda|consultant|1         |
|20190910  |A1     |Linda|b2c       |1         |
|20190916  |A1     |Linda|b2c       |2         |
|20190911  |B1     |John |b2c       |1         |
|20190916  |B1     |John |consultant|1         |
|20190917  |B1     |John |b2c       |2         |
|20190910  |B1     |Linda|b2c       |1         |
|20190910  |C1     |Tom  |b2c       |1         |
|20190915  |C1     |John |consultant|1         |
|20190916  |C1     |RJohn|consultant|1         |
|20190911  |C1     |John |b2c       |1         |
|20190916  |C1     |Linda|consultant|1         |
|20190917  |C1     |John |b2c       |2         |
+----------+-------+-----+----------+----------+

It's not exactly the same as your expected outcome, since it's ordered by Event_date and account.

Your expected output doesn't seem to be consistent. Please check the numbers again, especially for B1. Also RJohn in the input data.

You can do something that partition by Account, Name and Type. Then you can order by account followed by event_date.

from pyspark.sql.functions import *
from pyspark.sql.window import Window
df = spark.createDataFrame(
    [
        ("20190910", "A1", "Linda", "b2c"),
        ("20190911", "A1", "Tom", "consultant"),
        ("20190912", "A1", "John", "b2c"),
        ("20190913", "A1", "Tom", "consultant"),
        ("20190914", "A1", "Tom", "consultant"),
        ("20190915", "A1", "Linda", "consultant"),
        ("20190916", "A1", "Linda", "b2c"),
        ("20190917", "B1", "John", "b2c"),
        ("20190916", "B1", "John", "consultant"),
        ("20190910", "B1", "Linda", "b2c"),
        ("20190911", "B1", "John", "b2c"),
        ("20190915", "C1", "John", "consultant"),
        ("20190916", "C1", "Linda", "consultant"),
        ("20190917", "C1", "John", "b2c"),
        ("20190916", "C1", "John", "consultant"),
        ("20190910", "C1", "Tom", "b2c"),
        ("20190911", "C1", "John", "b2c"),
    ],
    ["Event_date", "account", "name", "type"],
)
w = Window.partitionBy("account", "name", "type").orderBy("Event_date")
df = df.withColumn("row_number", row_number().over(w)).orderBy("account","Event_date")

You will get the output as below : enter image description here

Related