Skip to content

Instantly share code, notes, and snippets.

@pietheinstrengholt
Last active December 27, 2021 15:53
Show Gist options
  • Select an option

  • Save pietheinstrengholt/ca2a0fb71dfe98cd1e8452ba3d544adc to your computer and use it in GitHub Desktop.

Select an option

Save pietheinstrengholt/ca2a0fb71dfe98cd1e8452ba3d544adc to your computer and use it in GitHub Desktop.
current_date = datetime.today().date()
print(current_date)
# Prepare for merge - Added effective and end date
df_source_new = dfDataChanged.withColumn('src_current', lit(True)).withColumn('src_effectiveDate', lit(current_date)).withColumn('src_endDate', lit(None))
df_source_new.show()
# FULL Merge, join on key column and also high date column to make only join to the latest records
df_merge = dfDataOriginal.join(df_source_new, (df_source_new.src_Id == dfDataOriginal.Id), how='fullouter')
# Derive new column to indicate the action
df_merge = df_merge.withColumn('action',
when(concat(df_merge.Id, df_merge.Firstname, df_merge.Lastname, df_merge.Department, df_merge.Salary) == concat(df_merge.src_Id, df_merge.src_Firstname, df_merge.src_Lastname, df_merge.src_Department, df_merge.src_Salary), 'NOACTION')
.when(df_merge.current == False, 'NOACTION')
.when(df_merge.src_Id.isNull() & df_merge.current, 'DELETE')
.when(df_merge.Id.isNull(), 'INSERT')
.otherwise('UPDATE')
)
df_merge.show()
# Generate the new data frames based on action code
column_names = ['Id', 'Firstname', 'Lastname', 'CreatedAt', 'Department', 'Salary', 'current', 'effectiveDate', 'endDate']
# For records that needs no action
df_merge_p1 = df_merge.filter(df_merge.action == 'NOACTION').select(column_names)
df_merge_p1.show()
# For records that needs insert only
df_merge_p2 = df_merge.filter(df_merge.action == 'INSERT').select(df_merge.src_Id.alias('id'),df_merge.src_Firstname.alias('Firstname'),df_merge.src_Lastname.alias('Lastname'),df_merge.src_CreatedAt.alias('CreatedAt'),df_merge.src_Department.alias('Department'),df_merge.src_Salary.alias('Salary'),lit(True).alias('current'),df_merge.src_effectiveDate.alias('effectiveDate'),df_merge.src_endDate.alias('endDate'))
df_merge_p2.show()
# For records that needs to be deleted
df_merge_p3 = df_merge.filter(df_merge.action == 'DELETE').select(column_names).withColumn('current', lit(False)).withColumn('endDate', lit(current_date))
df_merge_p3.show()
# For records that needs to be expired and then inserted
df_merge_p4_1 = df_merge.filter(df_merge.action == 'UPDATE').select(df_merge.src_Id.alias('Id'),df_merge.src_Firstname.alias('Firstname'),df_merge.src_Lastname.alias('Lastname'),df_merge.src_CreatedAt.alias('CreatedAt'),df_merge.src_Department.alias('Department'),df_merge.src_Salary.alias('Salary'),lit(True).alias('current'),df_merge.src_effectiveDate.alias('effectiveDate'),df_merge.src_endDate.alias('endDate'))
df_merge_p4_1.show()
df_merge_p4_2 = df_merge.filter(df_merge.action == 'UPDATE').withColumn('endDate', date_sub(df_merge.src_effectiveDate, 1)).withColumn('current', lit(False)).select(column_names)
df_merge_p4_2.show()
# Union all records together
df_merge_final = df_merge_p1.unionAll(df_merge_p2).unionAll(df_merge_p3).unionAll(df_merge_p4_1).unionAll(df_merge_p4_2)
df_merge_final.show()
# At last, you can overwrite existing data using this new data frame.
df_merge_final.write.format("delta").mode("overwrite").save(dfDataOriginalPath)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment