diff --git a/models/marts/reporting/rep_exchange__wallet_cash_flow_analysis.sql b/models/marts/reporting/rep_exchange__wallet_cash_flow_analysis.sql index 574812b..e3ab9d6 100644 --- a/models/marts/reporting/rep_exchange__wallet_cash_flow_analysis.sql +++ b/models/marts/reporting/rep_exchange__wallet_cash_flow_analysis.sql @@ -57,12 +57,15 @@ final__rep_exchange__wallet_cash_flow_analysis AS ( count_daily_wallet_activity, count_total_daily_transaction_credit, count_total_daily_transaction_debit, - prev_daily_last_wallet_balance, - daily_first_wallet_balance, - daily_last_wallet_balance, total_daily_transaction_credit, total_daily_transaction_debit, - + original_daily_first_wallet_balance, + original_daily_last_wallet_balance, + daily_0_origin_first_wallet_balance, + daily_0_origin_last_wallet_balance, + daily_calculated_first_wallet_balance, + daily_calculated_last_wallet_balance, + --date/time transaction_date, company_created_datetime, diff --git a/models/marts/warehouse/fct_exchange__wallet_daily_transactions.sql b/models/marts/warehouse/fct_exchange__wallet_daily_transactions.sql index 556e406..ebd9dbf 100644 --- a/models/marts/warehouse/fct_exchange__wallet_daily_transactions.sql +++ b/models/marts/warehouse/fct_exchange__wallet_daily_transactions.sql @@ -162,7 +162,7 @@ wallet_balance_list_join_wallet_transaction AS ( ), -data_cleaning AS ( +minor_data_imputation AS ( SELECT wallet_id, @@ -173,54 +173,154 @@ data_cleaning AS ( count_daily_wallet_activity, count_total_daily_transaction_credit, count_total_daily_transaction_debit, - - LAG(daily_last_wallet_balance) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_daily_last_wallet_balance, - + daily_last_wallet_balance, + CASE WHEN ROW_NUMBER() OVER (PARTITION BY wallet_id ORDER BY transaction_date DESC) = 1 THEN 1 ELSE 0 END AS is_latest_wallet_balance, - - CASE - /* - When opening balance does not equal to previous day's closing balance - Take previous day's closing balance as opening balance - */ - WHEN - daily_first_wallet_balance != prev_daily_last_wallet_balance - AND - daily_first_wallet_balance IS NOT NULL - THEN prev_daily_last_wallet_balance - /* - When there is transaction, but no wallet balance log - Take daily credit as first wallet balance - */ - WHEN daily_first_wallet_balance IS NULL - THEN total_daily_transaction_credit - - ELSE daily_first_wallet_balance - END AS daily_first_wallet_balance, CASE - WHEN daily_last_wallet_balance IS NULL - THEN ( total_daily_transaction_credit + total_daily_transaction_debit ) - ELSE daily_last_wallet_balance - END AS daily_last_wallet_balance, + WHEN ROW_NUMBER() OVER (PARTITION BY wallet_id ORDER BY transaction_date) = 1 + THEN 1 + ELSE 0 + END AS is_earliest_wallet_balance, + /* + Initialize for older records without wallet balance. (2 scenario) + 1. Set closing balance as 0 + 2. Set closing balance as difference between credit and debit + */ + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_first_wallet_balance, 0) + ELSE daily_first_wallet_balance + END daily_first_wallet_balance, + + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_last_wallet_balance, 0) + ELSE daily_last_wallet_balance + END daily_0_origin_last_wallet_balance, + + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_last_wallet_balance, total_daily_transaction_credit + total_daily_transaction_debit) + ELSE daily_last_wallet_balance + END daily_calculated_last_wallet_balance, + + + CASE + WHEN wallet_registration_datetime IS NULL + THEN ( + SELECT MIN(wallet_registration_datetime) + FROM wallet_balance_list_join_wallet_transaction t1 + WHERE t1.wallet_id = wallet_balance_list_join_wallet_transaction.wallet_id + ) + END AS new_wallet_registration_datetime, + + COALESCE(total_daily_transaction_credit, 0) AS total_daily_transaction_credit, + COALESCE(total_daily_transaction_debit, 0) AS total_daily_transaction_debit, + transaction_date, + LEAD(transaction_date) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS next_transaction_date, + LAG(transaction_date) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_transaction_date, + wallet_registration_datetime, + first_daily_transaction_start_datetime, + first_daily_transaction_completed_datetime, + last_daily_transaction_start_datetime, + last_daily_transaction_completed_datetime + + FROM wallet_balance_list_join_wallet_transaction + +), + +recursive_imputation AS ( + + SELECT + *, + + daily_first_wallet_balance AS new_daily_0_origin_first_wallet_balance, + daily_0_origin_last_wallet_balance AS new_daily_0_origin_last_wallet_balance, + + daily_first_wallet_balance AS new_daily_calculated_first_wallet_balance, + daily_calculated_last_wallet_balance AS new_daily_calculated_last_wallet_balance, + + FROM minor_data_imputation + + WHERE is_earliest_wallet_balance = 1 + + UNION ALL + + SELECT + t1.*, + + t2.new_daily_0_origin_last_wallet_balance AS new_daily_0_origin_first_wallet_balance, + t2.new_daily_0_origin_last_wallet_balance + t1.total_daily_transaction_credit + t1.total_daily_transaction_debit AS new_daily_0_origin_last_wallet_balance, + + t2.new_daily_calculated_last_wallet_balance AS new_daily_calculated_first_wallet_balance, + t2.new_daily_calculated_last_wallet_balance + t1.total_daily_transaction_credit + t1.total_daily_transaction_debit AS new_daily_calculated_last_wallet_balance + + + FROM minor_data_imputation t1 + + INNER JOIN recursive_imputation t2 + + ON ( + t1.wallet_id = t2.wallet_id + AND + t2.transaction_date = t1.prev_transaction_date + ) + +), + +casting_and_organize AS ( + + SELECT + wallet_id, + company_id, + currency_id, + first_wallet_log_id::STRING AS first_wallet_log_id, + last_wallet_log_id::STRING AS last_wallet_log_id, + is_earliest_wallet_balance, + is_latest_wallet_balance, + count_daily_wallet_activity, + count_total_daily_transaction_credit, + count_total_daily_transaction_debit, total_daily_transaction_credit, total_daily_transaction_debit, - wallet_registration_datetime, + LAG(daily_last_wallet_balance) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_daily_last_wallet_balance, + + COALESCE(CASE + WHEN prev_daily_last_wallet_balance IS NULL + THEN NULL + WHEN daily_first_wallet_balance != prev_daily_last_wallet_balance + THEN prev_daily_last_wallet_balance + ELSE daily_first_wallet_balance + END, 0) AS original_daily_first_wallet_balance, + + COALESCE(daily_last_wallet_balance, 0) AS original_daily_last_wallet_balance, + + COALESCE(new_daily_0_origin_first_wallet_balance, 0) AS daily_0_origin_first_wallet_balance, + COALESCE(new_daily_0_origin_last_wallet_balance, 0) AS daily_0_origin_last_wallet_balance, + + COALESCE(new_daily_calculated_first_wallet_balance, 0) AS daily_calculated_first_wallet_balance, + COALESCE(new_daily_calculated_last_wallet_balance, 0) AS daily_calculated_last_wallet_balance, + + prev_transaction_date, transaction_date, + next_transaction_date, first_daily_transaction_start_datetime, first_daily_transaction_completed_datetime, last_daily_transaction_start_datetime, last_daily_transaction_completed_datetime, - '{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime - - FROM wallet_balance_list_join_wallet_transaction + COALESCE(wallet_registration_datetime, new_wallet_registration_datetime) AS wallet_registration_datetime, + + '{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime + FROM recursive_imputation + ), @@ -236,30 +336,34 @@ final__fct_exchange__daily_wallet_transactions AS ( last_wallet_log_id, --dimensions + is_earliest_wallet_balance, is_latest_wallet_balance, --measures count_daily_wallet_activity, count_total_daily_transaction_credit, count_total_daily_transaction_debit, - prev_daily_last_wallet_balance, - daily_first_wallet_balance, - daily_last_wallet_balance, total_daily_transaction_credit, total_daily_transaction_debit, + original_daily_first_wallet_balance, + original_daily_last_wallet_balance, + daily_0_origin_first_wallet_balance, + daily_0_origin_last_wallet_balance, + daily_calculated_first_wallet_balance, + daily_calculated_last_wallet_balance, --date/times - wallet_registration_datetime, transaction_date, first_daily_transaction_start_datetime, first_daily_transaction_completed_datetime, last_daily_transaction_start_datetime, last_daily_transaction_completed_datetime, + wallet_registration_datetime, --metadata _dbt_ran_datetime - FROM data_cleaning + FROM casting_and_organize ) diff --git a/models/staging/exchange/stg_exchange__transaction_wallets.sql b/models/staging/exchange/stg_exchange__transaction_wallets.sql index 3e8cf3a..0271377 100644 --- a/models/staging/exchange/stg_exchange__transaction_wallets.sql +++ b/models/staging/exchange/stg_exchange__transaction_wallets.sql @@ -18,93 +18,356 @@ status AS ( -- LOGIC -transactions_join_transaction_types_payment_methods_currencies_status AS ( +-- wallet transactions +daily_transaction_log AS ( SELECT - transactions.transaction_id AS transaction_wallet_id, - IFF(transactions.owner_type = 'App\\Models\\Wallet', transactions.owner_id, null) AS wallet_id, - - COALESCE(transaction_types.name, transactions.transaction_type::string) AS transaction_type, - - transactions.receiver_user_id AS company_id, - - COALESCE(payment_methods.name, transactions.payment_method::string) AS payment_method, - - transactions.payment_reference, - transactions.bill_number, - - transactions.base_value, - transactions.quote_value, - - transactions.base_currency_id, - transactions.quote_currency_id, - - transactions.base_to_quote_currency_exchange_rate, - transactions.base_tax, - transactions.base_service_charge, - - transactions.expired_datetime, - - COALESCE(status.name, transactions.status::string) AS status, - - transactions.deleted_datetime, - transactions.created_datetime, - transactions.updated_datetime, - '{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime - - FROM transactions - - LEFT JOIN transaction_types - ON (transactions.transaction_type = transaction_types.id) - - LEFT JOIN payment_methods - ON (transactions.payment_method = payment_methods.id) - - LEFT JOIN status - ON (transactions.status = status.id) - - - WHERE - transactions.owner_type = 'App\\Models\\Wallet' -), - - --- FINAL -final__stg_exchange__transaction_wallets AS ( - - SELECT - -- ids - transaction_wallet_id, wallet_id, company_id, base_currency_id, quote_currency_id, + DATE(updated_datetime) AS transaction_date, + + MIN(created_datetime) AS first_daily_transaction_start_datetime, + MAX(created_datetime) AS last_daily_transaction_start_datetime, + + MIN(updated_datetime) AS first_daily_transaction_completed_datetime, + MAX(updated_datetime) AS last_daily_transaction_completed_datetime, + + SUM(CASE + WHEN transaction_type IN ('PAYMENT', 'DEBIT_NOTE') + THEN -( base_value ) + ELSE 0 + END) AS total_daily_transaction_debit, + + SUM(CASE + WHEN transaction_type IN ('TOP_UP', 'CREDIT_NOTE') + THEN base_value + ELSE 0 + END) AS total_daily_transaction_credit, - -- dimensions - transaction_type, - payment_method, - status, - payment_reference, - bill_number, + COUNT(CASE + WHEN transaction_type IN ('TOP_UP', 'CREDIT_NOTE') + THEN wallet_id + END) AS count_total_daily_transaction_credit, - -- measures - base_value, - quote_value, - base_to_quote_currency_exchange_rate, - base_tax, - base_service_charge, + COUNT(CASE + WHEN transaction_type IN ('PAYMENT', 'DEBIT_NOTE') + THEN wallet_id + END) AS count_total_daily_transaction_debit + + FROM wallet_transaction_logs - -- date/times - expired_datetime, - deleted_datetime, - created_datetime, + GROUP BY + wallet_id, + company_id, + base_currency_id, + quote_currency_id, + transaction_date + +), + +-- wallet balances +wallet_balance_list AS ( + + SELECT + wallet_log_id, + wallet_id, + company_id, + currency_id, + wallet_balance, + created_datetime AS wallet_registration_datetime, + updated_datetime + FROM wallet_balance_logs + + UNION + + SELECT + MD5_NUMBER_LOWER64(CONCAT(wallet_id, wallet_balance, updated_datetime)) AS wallet_log_id, + wallet_id, + company_id, + currency_id, + wallet_balance, + created_datetime AS wallet_registration_datetime, + updated_datetime + FROM latest_wallet_balances + +), + +wallet_balance_list_enrich AS ( + + SELECT + wallet_log_id, + wallet_id, + company_id, + currency_id, + wallet_balance, + wallet_registration_datetime, updated_datetime, + DATE(updated_datetime) AS updated_date, + LAG(wallet_balance) OVER (PARTITION BY wallet_id ORDER BY updated_datetime) AS prev_wallet_balance, + LEAD(wallet_balance) OVER (PARTITION BY wallet_id ORDER BY updated_datetime) AS next_wallet_balance + + FROM wallet_balance_list - -- metadata +), + +wallet_balance_daily_aggregate AS ( + + SELECT + DISTINCT + wallet_id, + company_id, + currency_id, + wallet_registration_datetime, + updated_date, + COUNT(wallet_id) OVER (PARTITION BY wallet_id, updated_date) AS count_daily_wallet_activity, + FIRST_VALUE(wallet_log_id) OVER (PARTITION BY wallet_id, updated_date ORDER BY updated_datetime) AS first_wallet_log_id, + LAST_VALUE(wallet_log_id) OVER (PARTITION BY wallet_id, updated_date ORDER BY updated_datetime) AS last_wallet_log_id, + + FIRST_VALUE(wallet_balance) OVER (PARTITION BY wallet_id, updated_date ORDER BY updated_datetime) AS daily_first_wallet_balance, + LAST_VALUE(wallet_balance) OVER (PARTITION BY wallet_id, updated_date ORDER BY updated_datetime) AS daily_last_wallet_balance, + + FROM wallet_balance_list_enrich + +), + +wallet_balance_list_join_wallet_transaction AS ( + + SELECT + COALESCE(wallet_balance_daily_aggregate.updated_date, daily_transaction_log.transaction_date) AS transaction_date, + COALESCE(wallet_balance_daily_aggregate.wallet_id, daily_transaction_log.wallet_id) AS wallet_id, + COALESCE(wallet_balance_daily_aggregate.company_id, daily_transaction_log.company_id) AS company_id, + COALESCE(wallet_balance_daily_aggregate.currency_id, daily_transaction_log.base_currency_id) AS currency_id, + + -- Wallet data + COALESCE(wallet_balance_daily_aggregate.count_daily_wallet_activity, 0) AS count_daily_wallet_activity, + COALESCE(wallet_balance_daily_aggregate.first_wallet_log_id, NULL) AS first_wallet_log_id, + COALESCE(wallet_balance_daily_aggregate.last_wallet_log_id, NULL) AS last_wallet_log_id, + COALESCE(wallet_balance_daily_aggregate.daily_first_wallet_balance, NULL) AS daily_first_wallet_balance, + COALESCE(wallet_balance_daily_aggregate.daily_last_wallet_balance, NULL) AS daily_last_wallet_balance, + COALESCE(wallet_balance_daily_aggregate.wallet_registration_datetime, NULL) AS wallet_registration_datetime, + + -- Transaction data + COALESCE(daily_transaction_log.first_daily_transaction_start_datetime, NULL) AS first_daily_transaction_start_datetime, + COALESCE(daily_transaction_log.last_daily_transaction_start_datetime, NULL) AS last_daily_transaction_start_datetime, + COALESCE(daily_transaction_log.first_daily_transaction_completed_datetime, NULL) AS first_daily_transaction_completed_datetime, + COALESCE(daily_transaction_log.last_daily_transaction_completed_datetime, NULL) AS last_daily_transaction_completed_datetime, + COALESCE(daily_transaction_log.total_daily_transaction_debit, NULL) AS total_daily_transaction_debit, + COALESCE(daily_transaction_log.total_daily_transaction_credit, NULL) AS total_daily_transaction_credit, + COALESCE(daily_transaction_log.count_total_daily_transaction_credit, 0) AS count_total_daily_transaction_credit, + COALESCE(daily_transaction_log.count_total_daily_transaction_debit, 0) AS count_total_daily_transaction_debit, + + FROM wallet_balance_daily_aggregate + + FULL OUTER JOIN daily_transaction_log + ON (wallet_balance_daily_aggregate.wallet_id = daily_transaction_log.wallet_id) + AND (wallet_balance_daily_aggregate.updated_date = daily_transaction_log.transaction_date) + +), + +minor_data_imputation AS ( + + SELECT + wallet_id, + company_id, + currency_id, + first_wallet_log_id, + last_wallet_log_id, + count_daily_wallet_activity, + count_total_daily_transaction_credit, + count_total_daily_transaction_debit, + daily_last_wallet_balance, + + CASE + WHEN ROW_NUMBER() OVER (PARTITION BY wallet_id ORDER BY transaction_date DESC) = 1 + THEN 1 + ELSE 0 + END AS is_latest_wallet_balance, + + CASE + WHEN ROW_NUMBER() OVER (PARTITION BY wallet_id ORDER BY transaction_date) = 1 + THEN 1 + ELSE 0 + END AS is_earliest_wallet_balance, + + /* + Initialize for older records without wallet balance. (2 scenario) + 1. Set closing balance as 0 + 2. Set closing balance as difference between credit and debit + */ + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_first_wallet_balance, 0) + ELSE daily_first_wallet_balance + END daily_first_wallet_balance, + + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_last_wallet_balance, 0) + ELSE daily_last_wallet_balance + END daily_0_origin_last_wallet_balance, + + CASE + WHEN is_earliest_wallet_balance = 1 + THEN COALESCE(daily_last_wallet_balance, total_daily_transaction_credit + total_daily_transaction_debit) + ELSE daily_last_wallet_balance + END daily_calculated_last_wallet_balance, + + + CASE + WHEN wallet_registration_datetime IS NULL + THEN ( + SELECT MIN(wallet_registration_datetime) + FROM wallet_balance_list_join_wallet_transaction t1 + WHERE t1.wallet_id = wallet_balance_list_join_wallet_transaction.wallet_id + ) + END AS new_wallet_registration_datetime, + + total_daily_transaction_credit, + total_daily_transaction_debit, + transaction_date, + LEAD(transaction_date) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS next_transaction_date, + LAG(transaction_date) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_transaction_date, + wallet_registration_datetime, + first_daily_transaction_start_datetime, + first_daily_transaction_completed_datetime, + last_daily_transaction_start_datetime, + last_daily_transaction_completed_datetime + + FROM wallet_balance_list_join_wallet_transaction + +), + +recursive_imputation AS ( + + SELECT + *, + + daily_first_wallet_balance AS new_daily_0_origin_first_wallet_balance, + daily_0_origin_last_wallet_balance AS new_daily_0_origin_last_wallet_balance, + + daily_first_wallet_balance AS new_daily_calculated_first_wallet_balance, + daily_calculated_last_wallet_balance AS new_daily_calculated_last_wallet_balance, + + FROM minor_data_imputation + + WHERE is_earliest_wallet_balance = 1 + + UNION ALL + + SELECT + t1.*, + + t2.new_daily_0_origin_last_wallet_balance AS new_daily_0_origin_first_wallet_balance, + t2.new_daily_0_origin_last_wallet_balance + t1.total_daily_transaction_credit + t1.total_daily_transaction_debit AS new_daily_0_origin_last_wallet_balance, + + t2.new_daily_calculated_last_wallet_balance AS new_daily_calculated_first_wallet_balance, + t2.new_daily_calculated_last_wallet_balance + t1.total_daily_transaction_credit + t1.total_daily_transaction_debit AS new_daily_calculated_last_wallet_balance + + + FROM minor_data_imputation t1 + + INNER JOIN recursive_imputation t2 + + ON ( + t1.wallet_id = t2.wallet_id + AND + t2.transaction_date = t1.prev_transaction_date + ) + +), + +casting_and_organize AS ( + + SELECT + wallet_id, + company_id, + currency_id, + first_wallet_log_id::STRING AS first_wallet_log_id, + last_wallet_log_id::STRING AS last_wallet_log_id, + is_earliest_wallet_balance + is_latest_wallet_balance, + count_daily_wallet_activity, + count_total_daily_transaction_credit, + count_total_daily_transaction_debit, + total_daily_transaction_credit, + total_daily_transaction_debit, + LAG(daily_last_wallet_balance) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_daily_last_wallet_balance, + + CASE + WHEN prev_daily_last_wallet_balance IS NULL + THEN NULL + WHEN daily_first_wallet_balance != prev_daily_last_wallet_balance + THEN prev_daily_last_wallet_balance + ELSE daily_first_wallet_balance + END AS original_daily_first_wallet_balance, + + daily_last_wallet_balance AS original_daily_last_wallet_balance, + + new_daily_0_origin_first_wallet_balance AS daily_0_origin_first_wallet_balance, + new_daily_0_origin_last_wallet_balance AS daily_0_origin_last_wallet_balance, + + new_daily_calculated_first_wallet_balance AS daily_calculated_first_wallet_balance, + new_daily_calculated_last_wallet_balance AS daily_calculated_last_wallet_balance, + + prev_transaction_date, + transaction_date, + next_transaction_date, + first_daily_transaction_start_datetime, + first_daily_transaction_completed_datetime, + last_daily_transaction_start_datetime, + last_daily_transaction_completed_datetime, + + COALESCE(wallet_registration_datetime, new_wallet_registration_datetime) AS wallet_registration_datetime, + + '{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime + + FROM recursive_imputation + +), + + +-- FINAL +final__fct_exchange__daily_wallet_transactions AS ( + + SELECT + --ids + wallet_id, + company_id, + currency_id, + first_wallet_log_id, + last_wallet_log_id, + + --dimensions + is_latest_wallet_balance, + + --measures + count_daily_wallet_activity, + count_total_daily_transaction_credit, + count_total_daily_transaction_debit, + total_daily_transaction_credit, + total_daily_transaction_debit, + original_daily_first_wallet_balance, + original_daily_last_wallet_balance, + daily_0_origin_first_wallet_balance, + daily_0_origin_last_wallet_balance, + daily_calculated_first_wallet_balance, + daily_calculated_last_wallet_balance, + + --date/times + transaction_date, + first_daily_transaction_start_datetime, + first_daily_transaction_completed_datetime, + last_daily_transaction_start_datetime, + last_daily_transaction_completed_datetime, + wallet_registration_datetime, + + --metadata _dbt_ran_datetime - FROM transactions_join_transaction_types_payment_methods_currencies_status + FROM casting_and_organize -) +) -SELECT * FROM final__stg_exchange__transaction_wallets \ No newline at end of file +SELECT * FROM final__fct_exchange__daily_wallet_transactions \ No newline at end of file