Last active
December 27, 2021 15:53
-
-
Save pietheinstrengholt/ca2a0fb71dfe98cd1e8452ba3d544adc to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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