I have this UPDATE SQL query that I need to convert to PySpark to work with dataframes. I'd like to know if it's possible to do it with dataframes and how to do it.
The SQL query:
UPDATE TBL1
SET COL_C=1
FROM TBL1
INNER JOIN TBL2 ON TBL1.COL_A=TBL2.COL_A AND TBL1.COL_B=TBL2.COL_B
INNER JOIN TBL3 ON TBL2.COL_A=TBL3.COL_A AND TBL2.COL_B=TBL3.COL_B
df_TBL1=TBL1
+-------+--------+----------+------+-----+
| COL_A| COL_B| dob|gender|COL_C|
+-------+--------+----------+------+-----+
| James| Smith|1991-04-01| M| 3000|
|Michael| Rose|2000-05-19| M| 4000|
| Robert|Williams|1978-09-05| M| 4000|
| Maria| Jones|1967-12-01| F| 4000|
| Jen| Brown|1980-02-17| F| 1000|
+-------+--------+----------+------+-----+
df_TBL2=TBL2
+-------+---------+----------+------+-----+
| COL_A| COL_B| dob|gender|COL_C|
+-------+---------+----------+------+-----+
| John| Snow|1791-04-01| M| 9000|
|Michael| Rose|2000-05-19| M| 4000|
| Robert|Baratheon|1778-09-05| M| 9500|
| Maria| Jones|1967-12-01| F| 4000|
+-------+---------+----------+------+-----+
df_TBL3=TBL3
+--------+------+----------+------+-----+
| COL_A| COL_B| dob|gender|COL_C|
+--------+------+----------+------+-----+
| Michael| Rose|2000-05-19| M| 4000|
| Peter|Parker|1978-09-05| M| 4000|
| Maria| Jones|1967-12-01| F| 4000|
|MaryJane| Brown|1980-12-17| F|10000|
+--------+------+----------+------+-----+
The joins give me:
df_TBL_ALL=df_TBL1 \
.join(df_TBL2,(df_TBL1.COL_A==df_TBL2.COL_A) & (df_TBL1.COL_B==df_TBL2.COL_B),how="inner") \
.join(df_TBL3,(df_TBL2.COL_A==df_TBL3.COL_A) & (df_TBL2.COL_B==df_TBL3.COL_B),how="inner") \
.select(df_TBL1["*"]) \
.withColumn("COL_C",spf.lit(1))
And then, I'm trying to join them
df_TBL1_JOINED=df_TBL1 \
.join(df_TBL_ALL,(df_TBL1.COL_A==df_TBL_ALL.COL_A) & (df_TBL1.COL_B==df_TBL_ALL.COL_B),how="left") \
.select(df_TBL1["*"], \
spf.coalesce(df_TBL_ALL.COL_C,df_TBL1.COL_C).alias("COL_C"))
df_TBL1_JOINED.show()
# +-------+--------+----------+------+-----+-----+
# | COL_A| COL_B| dob|gender|COL_C|COL_C|
# +-------+--------+----------+------+-----+-----+
# | James| Smith|1991-04-01| M| 3000| 3000|
# | Jen| Brown|1980-02-17| F| 1000| 1000|
# | Maria| Jones|1967-12-01| F| 4000| 1|
# |Michael| Rose|2000-05-19| M| 4000| 1|
# | Robert|Williams|1978-09-05| M| 4000| 4000|
# +-------+--------+----------+------+-----+-----+
But I'm confused about how to go on.
I did:
TBL01_R=TBL01_R \
.drop("COL_C")
TBL01_R=TBL01_R \
.withColumnRenamed("COL_Nova","COL_C").show()
TBL01=TBL01_R
# +-------+--------+----------+------+-----+
# | COL_A| COL_B| dob|gender|COL_C|
# +-------+--------+----------+------+-----+
# | James| Smith|1991-04-01| M| 3000|
# | Jen| Brown|1980-02-17| F| 1000|
# | Maria| Jones|1967-12-01| F| 1|
# |Michael| Rose|2000-05-19| M| 1|
# | Robert|Williams|1978-09-05| M| 4000|
# +-------+--------+----------+------+-----+
I got to the expected result but I don't know if it is the best performing way to achieve it.
Expected result: df_tbl1 with COL_C updated with a 1 in all rows present in the join of df_tbl1 with df_tbl2 and df_tbl3.
df_TBL1:
+-------+--------+----------+------+-----+
| COL_A| COL_B| dob|gender|COL_C|
+-------+--------+----------+------+-----+
| James| Smith|1991-04-01| M| 3000|
|Michael| Rose|2000-05-19| M| 1|
| Robert|Williams|1978-09-05| M| 4000|
| Maria| Jones|1967-12-01| F| 1|
| Jen| Brown|1980-02-17| F| 1000|
+-------+--------+----------+------+-----+