Merge branch 'main' of gitlab.com:cief-data/dbt_cloud into lee_build_exchange_rate_fact

This commit is contained in:
CIEF ACC1
2024-03-12 01:08:49 +00:00
12 changed files with 2446 additions and 43 deletions
@@ -0,0 +1,29 @@
version: 2
sources:
- name: src_c2m_exchange_rates
description: Currency exchange rate data of competitor C2M, obtained from web scraping. Link=https://www.c2m.my/
database: DEV_CIEF_RAW_DB
schema: C2M_WEBSCRAPE
loader: custom_python_code
loaded_at_field: SCRAPE_DATE
freshness:
warn_after: {count: 26, period: hour}
error_after: {count: 48, period: hour}
meta:
owner: "@yam"
model_maturity: prod
tags:
- exchange_rate_log
- daily
tables:
- name: currency_rate_logs
description: Web scraped daily currency exchange rates of different products offered by competitor C2M.
columns:
- name: service_type
description: (
ALTR:ALIPAY:Transfer/Scan-To-Pay,
TBDF:TAOBAO/1688:Payment-On-Behalf
)
@@ -0,0 +1,114 @@
-- IMPORT
WITH c2m_exchange_rate_log AS (
SELECT * FROM {{ source('src_c2m_exchange_rates', 'currency_rate_logs') }}
),
-- LOGIC
data_cleaning AS (
SELECT
*,
-- remove handling fee for TBDF
CASE
WHEN service_type = 'TBDF'
THEN 0
ELSE handeling_fee
END AS cleaned_handling_fee
FROM c2m_exchange_rate_log
),
rename_and_casting AS (
SELECT
product_name,
site_name AS competitor_name,
service_type,
pending_order AS order_queue_number,
myr::FLOAT AS base_value_myr,
cny::FLOAT AS quote_value_cny,
cleaned_handling_fee::FLOAT AS handling_fee_myr,
REPLACE(payable, ',','')::FLOAT AS payable_myr,
current_rate::FLOAT AS myr_to_cny_currency_exchange_rate,
staff_status,
DATEADD(DAY, -1, scrape_date) AS created_datetime,
scrape_hour,
TO_TIMESTAMP(
REGEXP_SUBSTR(
last_rate_update,
'[0-9]{4}-[0-9]{2}-[0-9]{2} [0-9]{1,2}:[0-9]{2} [APap][Mm]'
),
'YYYY-MM-DD HH12:MI AM'
) AS rate_updated_datetime
FROM data_cleaning
),
recalculate_payable_column AS (
SELECT
MD5_NUMBER_LOWER64(CONCAT(service_type, base_value_myr, created_datetime )) AS currency_rate_log_id,
product_name,
competitor_name,
service_type,
order_queue_number,
base_value_myr,
quote_value_cny,
handling_fee_myr,
myr_to_cny_currency_exchange_rate,
staff_status,
created_datetime,
scrape_hour,
rate_updated_datetime,
-- Recalculate payable myr after removing handling fee
CASE
WHEN service_type = 'TBDF'
THEN base_value_myr
ELSE payable_myr
END AS payable_myr,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM rename_and_casting
),
-- FINAL
final__stg_exchange__c2m_exchange_rates AS (
SELECT
-- ids
currency_rate_log_id,
-- dimensions
competitor_name,
product_name,
service_type,
staff_status,
-- measures
order_queue_number,
base_value_myr,
quote_value_cny,
myr_to_cny_currency_exchange_rate,
handling_fee_myr,
payable_myr,
-- date/time
rate_updated_datetime,
created_datetime,
-- metadata
_dbt_ran_datetime
FROM recalculate_payable_column
)
SELECT * FROM final__stg_exchange__c2m_exchange_rates
@@ -0,0 +1,29 @@
version: 2
sources:
- name: src_open_api_exchange_rates
description: Latest currency exchange rates of multiple currencies obtained from `Open Exchange Rates` API call.
database: DEV_CIEF_RAW_DB
schema: OPENEXCHANGERATES_AIRBYTE
loader: custom_python_code
loaded_at_field: _airbyte_extracted_at
freshness:
warn_after: {count: 26, period: hour}
error_after: {count: 48, period: hour}
meta:
owner: "@yam"
model_maturity: prod
tags:
- exchange_rate_log
- daily
tables:
- name: current_rate
description: Hourly currency exchange rates obtained from the API call.
columns:
- name: _airbyte_raw_id
description: Primary key for 'current_rate'
tests:
- unique
- not_null
@@ -0,0 +1,108 @@
-- IMPORT
WITH api_call_exchange_rates AS (
SELECT * FROM {{ source('src_open_api_exchange_rates', 'current_rate') }}
),
seed_historical_exchange_rates AS (
SELECT * FROM {{ ref('seed_exchange__historical_currency_exchange_rates') }}
),
-- LOGIC
api_rename_and_casting AS (
SELECT
_airbyte_raw_id AS _airbyte_id,
base AS base_currency,
PARSE_JSON(rates::STRING) AS parsed_rates,
_airbyte_extracted_at AS _airbyte_extracted_datetime,
TO_TIMESTAMP_NTZ(timestamp::INT) AS exchange_rate_datetime
FROM api_call_exchange_rates
),
api_currency_conversion AS (
SELECT
DATE(exchange_rate_datetime) AS exchange_rate_date,
AVG(1 / parsed_rates['MYR']) AS myr_to_usd_rate,
AVG(parsed_rates['CNY'] / parsed_rates['MYR']) AS myr_to_cny_rate
FROM api_rename_and_casting
GROUP BY
DATE(exchange_rate_datetime)
),
seed_currency_conversion AS (
SELECT
TO_CHAR(TO_DATE(date,'DD-MM-YYYY'), 'YYYY-MM-DD') AS exchange_rate_date,
cny / myr AS myr_to_cny_rate,
1 / myr AS myr_to_usd_rate
FROM seed_historical_exchange_rates
),
union_exchange_rates AS (
SELECT
exchange_rate_date,
myr_to_cny_rate,
myr_to_usd_rate
FROM seed_currency_conversion
WHERE
exchange_rate_date < '2024-01-04'
UNION ALL
SELECT
exchange_rate_date,
myr_to_cny_rate,
myr_to_usd_rate
FROM api_currency_conversion
),
add_primary_key_column AS (
SELECT
*,
MD5_NUMBER_LOWER64(CONCAT(myr_to_cny_rate, myr_to_usd_rate, exchange_rate_date)) AS currency_rate_log_id,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM union_exchange_rates
),
-- FINAL
final__stg_exchange__api_call_exchange_rates AS (
SELECT
-- id
currency_rate_log_id,
-- dimension
-- measure
myr_to_usd_rate,
myr_to_cny_rate,
-- date/time
exchange_rate_date,
-- metadata
_dbt_ran_datetime
FROM add_primary_key_column
)
SELECT * FROM final__stg_exchange__api_call_exchange_rates
@@ -41,9 +41,9 @@ currency_rate_logs_join_currency_rates AS (
LEFT JOIN currency_rates
ON (currency_rate_logs.currency_rate_id = currency_rates.currency_rate_id)
),
currency_rate_logs_quote_currency_id_join_currencies AS (
SELECT
@@ -58,7 +58,7 @@ currency_rate_logs_quote_currency_id_join_currencies AS (
currency_rate_logs_join_currency_rates.base_to_quote_currency_exchange_rate,
COALESCE(payment_methods.name, currency_rate_logs_join_currency_rates.payment_method::string) AS payment_method,
COALESCE(service_types.name, currency_rate_logs_join_currency_rates.service_type::string) AS service_type,
COALESCE(service_types.name, currency_rate_logs_join_currency_rates.service_type::string) AS service_type,
currency_rate_logs_join_currency_rates.creator_user_id,
currency_rate_logs_join_currency_rates.created_datetime,
@@ -75,66 +75,81 @@ currency_rate_logs_quote_currency_id_join_currencies AS (
LEFT JOIN service_types
ON (currency_rate_logs_join_currency_rates.service_type = service_types.service_type_id)
),
filtered_currency_rate_logs AS (
SELECT
*
FROM currency_rate_logs_quote_currency_id_join_currencies
WHERE service_type NOT IN ('x1', 'Company to Company', 'RECEIVED-ON-BEHALF')
),
currency_rate_logs_add_dense_rank AS (
SELECT
currency_rate_logs_quote_currency_id_join_currencies.currency_rate_log_id,
filtered_currency_rate_logs.currency_rate_log_id,
dense_rank() OVER
(PARTITION BY currency_rate_logs_quote_currency_id_join_currencies.currency_rate_id
ORDER BY currency_rate_logs_quote_currency_id_join_currencies.created_datetime) AS dense_rank_index,
DENSE_RANK() OVER
(PARTITION BY filtered_currency_rate_logs.currency_rate_id
ORDER BY filtered_currency_rate_logs.created_datetime) AS dense_rank_index,
currency_rate_logs_quote_currency_id_join_currencies.currency_rate_id AS currency_rate_id,
currency_rate_logs_quote_currency_id_join_currencies.quote_currency_id,
currency_rate_logs_quote_currency_id_join_currencies.quote_currency_name,
currency_rate_logs_quote_currency_id_join_currencies.quote_currency_short_code,
currency_rate_logs_quote_currency_id_join_currencies.quote_currency_symbol,
currency_rate_logs_quote_currency_id_join_currencies.payment_method,
currency_rate_logs_quote_currency_id_join_currencies.service_type,
currency_rate_logs_quote_currency_id_join_currencies.base_to_quote_currency_exchange_rate,
currency_rate_logs_quote_currency_id_join_currencies.creator_user_id,
currency_rate_logs_quote_currency_id_join_currencies.created_datetime
filtered_currency_rate_logs.currency_rate_id AS currency_rate_id,
filtered_currency_rate_logs.quote_currency_id,
filtered_currency_rate_logs.quote_currency_name,
filtered_currency_rate_logs.quote_currency_short_code,
filtered_currency_rate_logs.quote_currency_symbol,
filtered_currency_rate_logs.payment_method,
filtered_currency_rate_logs.service_type,
filtered_currency_rate_logs.base_to_quote_currency_exchange_rate,
filtered_currency_rate_logs.creator_user_id,
filtered_currency_rate_logs.created_datetime
FROM currency_rate_logs_quote_currency_id_join_currencies
FROM filtered_currency_rate_logs
ORDER BY
currency_rate_id,
created_datetime
),
currency_rate_logs_resturucure AS (
currency_rate_logs_resturucure AS (
SELECT
currency_rate_logs_add_dense_rank.currency_rate_log_id,
currency_rate_logs_add_dense_rank.currency_rate_id,
currency_rate_logs_add_dense_rank.quote_currency_id,
currency_rate_logs_add_dense_rank.quote_currency_name,
currency_rate_logs_add_dense_rank.quote_currency_short_code,
currency_rate_logs_add_dense_rank.quote_currency_symbol,
currency_rate_logs_add_dense_rank.payment_method,
currency_rate_logs_add_dense_rank.service_type,
currency_rate_logs_add_dense_rank.base_to_quote_currency_exchange_rate,
currency_rate_logs_add_dense_rank.creator_user_id,
SELECT
currency_rate_logs_add_dense_rank.currency_rate_log_id,
currency_rate_logs_add_dense_rank.currency_rate_id,
currency_rate_logs_add_dense_rank.quote_currency_id,
currency_rate_logs_add_dense_rank.quote_currency_name,
currency_rate_logs_add_dense_rank.quote_currency_short_code,
currency_rate_logs_add_dense_rank.quote_currency_symbol,
currency_rate_logs_add_dense_rank.payment_method,
currency_rate_logs_add_dense_rank.service_type,
currency_rate_logs_add_dense_rank.base_to_quote_currency_exchange_rate,
currency_rate_logs_add_dense_rank.creator_user_id,
CASE
WHEN currency_rate_logs_add_dense_rank.dense_rank_index=1 THEN CAST('1990-01-01 00:00:00 +08:00' AS DATETIME)
ELSE currency_rate_logs_add_dense_rank.created_datetime
END AS from_datetime,
CASE
WHEN currency_rate_logs_add_dense_rank.dense_rank_index=1 THEN CAST('1990-01-01 00:00:00 +08:00' AS DATETIME)
ELSE currency_rate_logs_add_dense_rank.created_datetime
END AS from_datetime,
COALESCE(currency_rate_logs_add_dense_rank_clone.created_datetime,
CAST('2999-12-31 23:59:59 +08:00' AS DATETIME)) AS to_datetime,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
COALESCE(currency_rate_logs_add_dense_rank_clone.created_datetime,
CAST('2999-12-31 23:59:59 +08:00' AS DATETIME)) AS to_datetime,
FROM
currency_rate_logs_add_dense_rank
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
LEFT JOIN currency_rate_logs_add_dense_rank AS currency_rate_logs_add_dense_rank_clone
ON currency_rate_logs_add_dense_rank.dense_rank_index = currency_rate_logs_add_dense_rank_clone.dense_rank_index-1
AND currency_rate_logs_add_dense_rank.currency_rate_id = currency_rate_logs_add_dense_rank_clone.currency_rate_id
FROM
currency_rate_logs_add_dense_rank
LEFT JOIN currency_rate_logs_add_dense_rank AS currency_rate_logs_add_dense_rank_clone
ON currency_rate_logs_add_dense_rank.dense_rank_index = currency_rate_logs_add_dense_rank_clone.dense_rank_index-1
AND currency_rate_logs_add_dense_rank.currency_rate_id = currency_rate_logs_add_dense_rank_clone.currency_rate_id
ORDER BY currency_rate_logs_add_dense_rank.currency_rate_id, currency_rate_logs_add_dense_rank.created_datetime
ORDER BY currency_rate_logs_add_dense_rank.currency_rate_id, currency_rate_logs_add_dense_rank.created_datetime
),
@@ -44,6 +44,16 @@ sources:
columns:
- name: _airbyte_raw_id
description: Primary key for 'account_frozen_case'
tests:
- unique
- not_null
- name: sd_currency_exchange_rates
description: Data retrieved from gsheet, comprises the currency exchange rates from our competitor SD.
columns:
- name: _airbyte_raw_id
description: Primary key for 'sd_currency_exchange_rates'
tests:
- unique
- not_null
@@ -0,0 +1,69 @@
-- IMPORT
WITH sd_exchange_rate_log AS (
SELECT * FROM {{ source('src_googlesheet_airbyte', 'sd_currency_exchange_rates') }}
),
seed_metadata AS (
SELECT * FROM {{ ref('seed_exchange__metadata') }}
),
-- LOGIC
rename_and_casting AS (
SELECT
{{ dbt_utils.generate_surrogate_key(['payment_method', 'myr_to_cny_rate', 'exchange_rate_date']) }} AS sd_exchange_rate_id,
_airbyte_raw_id AS _airbyte_id,
sd_exchange_rate_log.provider AS competitor_name,
sd_exchange_rate_log.payment_method,
sd_exchange_rate_log.myr_to_cny_rate::FLOAT AS myr_to_cny_currency_exchange_rate,
sd_exchange_rate_log.transfer_days AS transfer_duration_remark,
sd_exchange_rate_log.remark,
COALESCE(
TRY_TO_DATE(sd_exchange_rate_log.exchange_rate_date, 'DD-MM-YY'),
TRY_TO_DATE(sd_exchange_rate_log.exchange_rate_date, 'MM-DD-YY')
) AS exchange_rate_date,
seed_metadata.value AS handling_fee_cny,
'{{ modules.datetime.datetime.now(modules.pytz.timezone("Asia/Kuala_Lumpur")) }}' AS _dbt_ran_datetime
FROM sd_exchange_rate_log
LEFT JOIN seed_metadata
ON (seed_metadata.table_name = 'stg_exchange__sd_currency_rate_logs')
),
-- FINAL
final__stg_exchange__sd_currency_rate_logs AS (
SELECT
-- ids
sd_exchange_rate_id,
-- dimensions
competitor_name,
payment_method,
transfer_duration_remark,
remark,
-- measures
handling_fee_cny,
myr_to_cny_currency_exchange_rate,
-- date/time
exchange_rate_date,
-- metadata
_dbt_ran_datetime
FROM rename_and_casting
)
SELECT * FROM final__stg_exchange__sd_currency_rate_logs