Workflow saved

This commit is contained in:
unknown
2025-09-11 09:20:44 +00:00
parent 2e42932fab
commit 2f13af0751
3 changed files with 19 additions and 17 deletions

View File

@@ -354,16 +354,7 @@ def failed_payments_mapper(failed_payments_reader_df, job_id, spark):
failed_payments_mapper_df=spark.sql(("SELECT " + ', '.join(_failed_payments_mapper_select_clause) + " FROM failed_payments_reader_df").replace("{job_id}",f"'{job_id}'"))
failed_payments_mapper_df.createOrReplaceTempView("failed_payments_mapper_df")
return (failed_payments_mapper_df,)
@app.cell
def final_failed_payments(failed_payments_mapper_df, spark):
print(failed_payments_mapper_df.columns)
final_failed_payments_df = spark.sql("select * from failed_payments_mapper_df where payment_date >= COALESCE((SELECT MAX(DATE(payment_date)) FROM dremio.failedpaymentmetrics), (SELECT MIN(payment_date) FROM failed_payments_mapper_df))")
final_failed_payments_df.createOrReplaceTempView('final_failed_payments_df')
return (final_failed_payments_df,)
return
@app.cell
@@ -578,5 +569,15 @@ def success_payment_metrics_writer(spark, success_payment_metrics_df):
return
@app.cell
def failed_payments(FailedPaymentsData_df, failed_payments_df, spark):
print(FailedPaymentsData_df.columns)
LatestFailedPayments_df = spark.sql("select * from FailedPaymentsData_df where payment_date >= COALESCE((SELECT MAX(DATE(payment_date)) FROM dremio.failedpaymentmetrics), (SELECT MIN(payment_date) FROM failed_payments_mapper_df))")
failed_payments_df.createOrReplaceTempView('failed_payments_df')
failed_payments_df.persist()
return
if __name__ == "__main__":
app.run()