cuteabhi32
New Contributor III

from pyspark.sql.functions import * 

from pyspark.sql.functions import col

import pyspark.sql.functions as F

from datetime import date,datetime

import time 

from dateutil.relativedelta import relativedelta

from dateutil import parser

from pyspark.sql.window import Window

from pyspark.sql.types import *

import locale

# ---- merge statements ---- 

df1 = spark.read.format("csv").option("header", "true").load("dbfs:/FileStore/shared_uploads/a.csv")

df1.show()

df2 = spark.read.format("csv").option("header", "true").load("dbfs:/FileStore/shared_uploads/b.csv")

df2.show()

# ---- merge statements ---- 

sf_start_dt = '02MAY2022'

a=df1

a = a.filter(col("ref_date") == f"{sf_start_dt}")

a = a.drop('name')

a = a.select('id','salary','ref_date','std','mean',)

b = df2

a.createOrReplaceTempView("a")

b.createOrReplaceTempView("b")

a = a.join(b,'id',"outer")

df1 = spark.sql("select tbl1.id,(select 1) as tempCol1 from a tbl1 inner join b tbl2 on tbl1.id = tbl2.id")

df2 = spark.sql("select tbl1.id,(select 2) as tempCol2 from a tbl1 left join b tbl2 on tbl1.id = tbl2.id where tbl2.id is null")

df3 = spark.sql("select tbl1.id,(select 3) as tempCol3 from b tbl1 left join a tbl2 on tbl1.id = tbl2.id where tbl2.id is null")

a = a.join(df1,'id',"outer").join(df2,'id',"outer").join(df3,'id',"outer")

a = a.na.fill(0,'tempCol1')

a = a.na.fill(0,'tempCol2')

a = a.na.fill(0,'tempCol3')

a = a.withColumn('flag', coalesce(col('tempCol1')+col('tempCol2')+col('tempCol3')) )

a = a.drop('tempCol1')

a = a.drop('tempCol2')

a = a.drop('tempCol3')

a = a\

.withColumn("neg_std",F.expr(f"(std*(-1))"))

a = a\

.withColumn("mean20perc",F.expr(f"(0.20*mean)"))

a = a\

.withColumn("neg_mean20perc",F.expr(f"(mean20perc*(-1))"))

a = a\

.withColumn("new_var",F.expr(f"'{sf_start_dt}'"))

columnsToDrop = []

selectClause = ''

a.createOrReplaceTempView("a")

a = spark.sql("select * from a")

from pyspark.sql.types import StringType

from pyspark.sql.functions import udf

column =list(a.columns)

print(column)

def func_udf(df,col):

  column =list(df.columns)

  if col in column:

    return df

  else:

     df.withColumn("col", lit(null))

spark.udf.register("ColumnChecker",func_udf)

a = a.withColumn('ref_date',expr(f"CASE WHEN flag = 3 THEN '{sf_start_dt}' ELSE ColumnChecker(a,ref_date) END"))

a = a.withColumn('balance',expr(f"CASE WHEN flag = 3 THEN 0 END"))

a.show()

a = a.withColumn('new_col',expr(f"CASE WHEN flag = 3 THEN '{sf_start_dt}' END"))

a.show()

work_ppcin_bal2_2019_1 = a

work_ppcin_bal2_2019_1.show()

# ---- end of merge statements ---- 

this is the full fledge code