- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
06-07-2022 03:58 AM
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