build wallet cash flow analysis

This commit is contained in:
CIEF ACC1
2024-03-15 07:15:38 +00:00
parent 05339a85ad
commit 0550be2ed3
2 changed files with 348 additions and 0 deletions
@@ -0,0 +1,82 @@
-- IMPORT
WITH daily_wallet_transactions AS (
SELECT * FROM {{ ref('fct_exchange__wallet_daily_transactions') }}
),
companies AS (
SELECT * FROM {{ ref('dim_exchange__companies') }}
),
-- LOGIC
daily_wallet_transaction_enrich AS (
SELECT
daily_wallet_transactions.* EXCLUDE _dbt_ran_datetime,
companies.company_marking_id,
companies.company_type,
companies.business_type,
companies.exchange_rate_segment,
companies.is_migrated_company,
companies.company_created_datetime,
MD5_NUMBER_LOWER64(CONCAT(daily_wallet_transactions.wallet_id, transaction_date))::STRING AS wallet_cash_flow_transaction_id,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM daily_wallet_transactions
LEFT JOIN companies
ON ( daily_wallet_transactions.company_id = companies.company_id)
),
-- FINAL
final__rep_exchange__wallet_cash_flow_analysis AS (
SELECT
--ids
wallet_cash_flow_transaction_id,
wallet_id,
company_id,
company_marking_id,
currency_id,
first_wallet_log_id,
last_wallet_log_id,
--dimensions
company_type,
business_type,
exchange_rate_segment,
is_migrated_company,
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,
--date/time
transaction_date,
company_created_datetime,
wallet_registration_datetime,
first_daily_transaction_start_datetime,
first_daily_transaction_completed_datetime,
last_daily_transaction_start_datetime,
last_daily_transaction_completed_datetime,
--metadata
_dbt_ran_datetime
FROM daily_wallet_transaction_enrich
)
SELECT * FROM final__rep_exchange__wallet_cash_flow_analysis
@@ -0,0 +1,266 @@
-- IMPORT
WITH wallet_balance_logs AS (
SELECT * FROM {{ ref('stg_exchange__wallet_balance_logs') }}
),
latest_wallet_balances AS (
SELECT * FROM {{ ref('stg_exchange__wallet_balances') }}
),
wallet_transaction_logs AS (
SELECT * FROM {{ ref('stg_exchange__transaction_wallet_logs') }}
WHERE status IN ('APPROVED')
),
-- LOGIC
-- wallet transactions
daily_transaction_log AS (
SELECT
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,
COUNT(CASE
WHEN transaction_type IN ('TOP_UP', 'CREDIT_NOTE')
THEN wallet_id
END) AS count_total_daily_transaction_credit,
COUNT(CASE
WHEN transaction_type IN ('PAYMENT', 'DEBIT_NOTE')
THEN wallet_id
END) AS count_total_daily_transaction_debit
FROM wallet_transaction_logs
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
),
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)
),
data_cleaning 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,
LAG(daily_last_wallet_balance) OVER (PARTITION BY wallet_id ORDER BY transaction_date) AS prev_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,
total_daily_transaction_credit,
total_daily_transaction_debit,
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,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM wallet_balance_list_join_wallet_transaction
),
-- 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,
prev_daily_last_wallet_balance,
daily_first_wallet_balance,
daily_last_wallet_balance,
total_daily_transaction_credit,
total_daily_transaction_debit,
--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,
--metadata
_dbt_ran_datetime
FROM data_cleaning
)
SELECT * FROM final__fct_exchange__daily_wallet_transactions