Merge branch 'build-wallet-cash-flow-analysis' into 'main'

Build wallet cash flow analysis

See merge request cief-data/dbt_cloud!36
This commit is contained in:
CIEF ACC1
2024-05-08 08:27:33 +00:00
7 changed files with 690 additions and 6 deletions
@@ -0,0 +1,85 @@
-- 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,
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,
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,370 @@
-- 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)
),
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,
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,
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,
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_earliest_wallet_balance,
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 casting_and_organize
)
SELECT * FROM final__fct_exchange__daily_wallet_transactions
@@ -0,0 +1,58 @@
-- IMPORT
WITH source AS (
SELECT * FROM {{ source('src_exchange_mysql', 'wallet_logs') }}
),
-- LOGIC
renamed AS (
SELECT
id AS wallet_log_id,
wallet_id,
owner_type,
owner_id,
code,
currency_id,
amount,
deleted_at AS deleted_datetime,
created_at AS created_datetime,
updated_at AS updated_datetime,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM source
),
-- FINAL
final__base_exchange__wallet_logs AS (
SELECT
-- ids
wallet_log_id,
wallet_id,
owner_id,
currency_id,
-- dimensions
code,
owner_type,
-- measures
amount,
-- date/time
created_datetime,
updated_datetime,
deleted_datetime,
-- metadata
_dbt_ran_datetime
FROM renamed
)
SELECT * FROM final__base_exchange__wallet_logs
@@ -0,0 +1,56 @@
-- IMPORT
WITH source AS (
SELECT * FROM {{ source('src_exchange_mysql', 'wallets') }}
),
-- LOGIC
renamed AS (
SELECT
id AS wallet_id,
owner_type,
owner_id,
code,
currency_id,
amount,
deleted_at AS deleted_datetime,
created_at AS created_datetime,
updated_at AS updated_datetime,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM source
),
-- FINAL
final__base_exchange__wallets AS (
SELECT
-- ids
wallet_id,
owner_id,
currency_id,
-- dimensions
code,
owner_type,
-- measures
amount,
-- date/time
created_datetime,
updated_datetime,
deleted_datetime,
-- metadata
_dbt_ran_datetime
FROM renamed
)
SELECT * FROM final__base_exchange__wallets
@@ -50,6 +50,7 @@ transactions_join_transaction_types_payment_methods_currencies_status AS (
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
@@ -63,9 +64,9 @@ transactions_join_transaction_types_payment_methods_currencies_status AS (
LEFT JOIN status
ON (transactions.status = status.id)
WHERE
transactions.owner_type = 'App\\Models\\Wallet'
),
@@ -0,0 +1,58 @@
-- IMPORT
WITH wallet_logs AS (
SELECT * FROM {{ ref('base_exchange__wallet_logs') }}
),
-- LOGIC
renamed AS (
SELECT
wallet_log_id,
wallet_id,
owner_id AS company_id,
currency_id,
code,
amount AS wallet_balance,
created_datetime,
updated_datetime,
deleted_datetime,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM wallet_logs
WHERE owner_type = 'App\\Models\\Company'
),
-- FINAL
final__base_exchange__wallet_logs AS (
SELECT
--ids
wallet_log_id,
wallet_id,
company_id,
currency_id,
--dimensions
code,
--measures
wallet_balance,
--date/times
created_datetime,
updated_datetime,
deleted_datetime,
--metadata
_dbt_ran_datetime
FROM renamed
)
SELECT * FROM final__base_exchange__wallet_logs
@@ -0,0 +1,56 @@
-- IMPORT
WITH wallets AS (
SELECT * FROM {{ ref('base_exchange__wallets') }}
),
-- LOGIC
renamed AS (
SELECT
wallet_id,
owner_id AS company_id,
currency_id,
code,
amount AS wallet_balance,
created_datetime,
updated_datetime,
deleted_datetime,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM wallets
WHERE owner_type = 'App\\Models\\Company'
),
-- FINAL
final__stg__exchange__wallets AS (
SELECT
--ids
wallet_id,
company_id,
currency_id,
--dimensions
code,
--measures
wallet_balance,
--date/times
created_datetime,
updated_datetime,
deleted_datetime,
--metadata
_dbt_ran_datetime
FROM renamed
)
SELECT * FROM final__stg__exchange__wallets