Skip to content

Commit ccdd676

Browse files
update ingestion for commodities
1 parent 4eb09a5 commit ccdd676

14 files changed

Lines changed: 1276 additions & 28 deletions
Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,112 @@
1+
{{
2+
config(
3+
materialized='table',
4+
description='Agricultural commodities daily data aggregated to monthly averages with forward-looking quarterly percentage changes'
5+
)
6+
}}
7+
8+
WITH monthly_averages AS (
9+
SELECT
10+
commodity_name,
11+
commodity_unit,
12+
EXTRACT(YEAR FROM date) AS year_val,
13+
EXTRACT(MONTH FROM date) AS month_val,
14+
CONCAT(
15+
EXTRACT(YEAR FROM date),
16+
'-',
17+
LPAD(EXTRACT(MONTH FROM date)::VARCHAR, 2, '0')
18+
) AS year_month,
19+
MAKE_DATE(EXTRACT(YEAR FROM date), EXTRACT(MONTH FROM date), 1)
20+
AS month_date,
21+
ROUND(AVG(price), 4) AS avg_price
22+
FROM {{ ref('stg_agriculture_commodities') }}
23+
GROUP BY
24+
commodity_name,
25+
commodity_unit,
26+
EXTRACT(YEAR FROM date),
27+
EXTRACT(MONTH FROM date)
28+
),
29+
30+
quarterly_data AS (
31+
SELECT
32+
commodity_name,
33+
commodity_unit,
34+
year_val,
35+
year_month,
36+
month_date,
37+
avg_price,
38+
CASE
39+
WHEN month_val IN (1, 2, 3) THEN 1
40+
WHEN month_val IN (4, 5, 6) THEN 2
41+
WHEN month_val IN (7, 8, 9) THEN 3
42+
WHEN month_val IN (10, 11, 12) THEN 4
43+
END AS quarter_num,
44+
CONCAT(
45+
year_val, '-Q',
46+
CASE
47+
WHEN month_val IN (1, 2, 3) THEN 1
48+
WHEN month_val IN (4, 5, 6) THEN 2
49+
WHEN month_val IN (7, 8, 9) THEN 3
50+
WHEN month_val IN (10, 11, 12) THEN 4
51+
END
52+
) AS year_quarter,
53+
AVG(avg_price) OVER (
54+
PARTITION BY
55+
commodity_name, commodity_unit, year_val,
56+
CASE
57+
WHEN month_val IN (1, 2, 3) THEN 1
58+
WHEN month_val IN (4, 5, 6) THEN 2
59+
WHEN month_val IN (7, 8, 9) THEN 3
60+
WHEN month_val IN (10, 11, 12) THEN 4
61+
END
62+
) AS quarterly_avg_price
63+
FROM monthly_averages
64+
),
65+
66+
with_forward_quarters AS (
67+
SELECT
68+
*,
69+
LEAD(quarterly_avg_price, 1) OVER (
70+
PARTITION BY commodity_name, commodity_unit
71+
ORDER BY year_val, quarter_num
72+
) AS price_q1_forward,
73+
LEAD(quarterly_avg_price, 2) OVER (
74+
PARTITION BY commodity_name, commodity_unit
75+
ORDER BY year_val, quarter_num
76+
) AS price_q2_forward,
77+
LEAD(quarterly_avg_price, 3) OVER (
78+
PARTITION BY commodity_name, commodity_unit
79+
ORDER BY year_val, quarter_num
80+
) AS price_q3_forward,
81+
LEAD(quarterly_avg_price, 4) OVER (
82+
PARTITION BY commodity_name, commodity_unit
83+
ORDER BY year_val, quarter_num
84+
) AS price_q4_forward
85+
FROM quarterly_data
86+
)
87+
88+
SELECT
89+
commodity_name,
90+
commodity_unit,
91+
year_month,
92+
month_date,
93+
year_quarter,
94+
quarter_num,
95+
year_val,
96+
avg_price AS monthly_avg_price,
97+
ROUND(quarterly_avg_price, 4) AS quarterly_avg_price,
98+
ROUND(
99+
(price_q1_forward - quarterly_avg_price) / quarterly_avg_price * 100, 2
100+
) AS pct_change_q1_forward,
101+
ROUND(
102+
(price_q2_forward - quarterly_avg_price) / quarterly_avg_price * 100, 2
103+
) AS pct_change_q2_forward,
104+
ROUND(
105+
(price_q3_forward - quarterly_avg_price) / quarterly_avg_price * 100, 2
106+
) AS pct_change_q3_forward,
107+
ROUND(
108+
(price_q4_forward - quarterly_avg_price) / quarterly_avg_price * 100, 2
109+
) AS pct_change_q4_forward
110+
FROM with_forward_quarters
111+
ORDER BY commodity_name, commodity_unit, month_date
112+
Lines changed: 191 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,191 @@
1+
{{ config(
2+
materialized='view'
3+
) }}
4+
5+
WITH base_data AS (
6+
SELECT
7+
commodity_name,
8+
commodity_unit,
9+
price,
10+
CAST(date AS DATE) AS trade_date,
11+
-- Calculate day-over-day price changes
12+
price - LAG(price) OVER (
13+
PARTITION BY commodity_name
14+
ORDER BY date
15+
) AS price_change,
16+
-- Calculate percentage changes
17+
CASE
18+
WHEN LAG(price) OVER (
19+
PARTITION BY commodity_name
20+
ORDER BY date
21+
) > 0
22+
THEN (
23+
(price - LAG(price) OVER (
24+
PARTITION BY commodity_name
25+
ORDER BY date
26+
))
27+
/ LAG(price) OVER (
28+
PARTITION BY commodity_name
29+
ORDER BY date
30+
)
31+
)
32+
* 100
33+
END AS pct_change
34+
FROM {{ ref('stg_agriculture_commodities') }}
35+
WHERE price IS NOT NULL
36+
AND date IS NOT NULL
37+
AND price > 0
38+
),
39+
40+
-- Define date boundaries for different periods
41+
date_boundaries AS (
42+
SELECT
43+
CURRENT_DATE AS today,
44+
CURRENT_DATE - INTERVAL '12 weeks' AS twelve_weeks_ago,
45+
CURRENT_DATE - INTERVAL '6 months' AS six_months_ago,
46+
CURRENT_DATE - INTERVAL '1 year' AS one_year_ago,
47+
CURRENT_DATE - INTERVAL '5 years' AS five_years_ago
48+
),
49+
50+
-- Filter data for each time period
51+
filtered_data AS (
52+
SELECT
53+
bd.*,
54+
CASE
55+
WHEN bd.trade_date >= db.twelve_weeks_ago THEN '12_weeks'
56+
WHEN bd.trade_date >= db.six_months_ago THEN '6_months'
57+
WHEN bd.trade_date >= db.one_year_ago THEN '1_year'
58+
WHEN bd.trade_date >= db.five_years_ago THEN '5_years'
59+
ELSE 'older'
60+
END AS time_period
61+
FROM base_data AS bd
62+
CROSS JOIN date_boundaries AS db
63+
WHERE bd.trade_date >= db.five_years_ago
64+
AND bd.price_change IS NOT NULL
65+
),
66+
67+
-- Get first and last prices for each period
68+
period_boundaries AS (
69+
SELECT
70+
commodity_name,
71+
commodity_unit,
72+
time_period,
73+
MIN(trade_date) AS period_start_date,
74+
MAX(trade_date) AS period_end_date
75+
FROM filtered_data
76+
WHERE time_period != 'older'
77+
GROUP BY commodity_name, commodity_unit, time_period
78+
),
79+
80+
-- Get start and end prices
81+
start_prices AS (
82+
SELECT
83+
pb.commodity_name,
84+
pb.commodity_unit,
85+
pb.time_period,
86+
fd.price AS period_start_price
87+
FROM period_boundaries AS pb
88+
INNER JOIN filtered_data AS fd ON
89+
pb.commodity_name = fd.commodity_name
90+
AND pb.time_period = fd.time_period
91+
AND pb.period_start_date = fd.trade_date
92+
),
93+
94+
end_prices AS (
95+
SELECT
96+
pb.commodity_name,
97+
pb.commodity_unit,
98+
pb.time_period,
99+
fd.price AS period_end_price
100+
FROM period_boundaries AS pb
101+
INNER JOIN filtered_data AS fd ON
102+
pb.commodity_name = fd.commodity_name
103+
AND pb.time_period = fd.time_period
104+
AND pb.period_end_date = fd.trade_date
105+
),
106+
107+
-- Main aggregation
108+
aggregated_results AS (
109+
SELECT
110+
commodity_name,
111+
commodity_unit,
112+
time_period,
113+
MIN(trade_date) AS period_start_date,
114+
MAX(trade_date) AS period_end_date,
115+
COUNT(*) AS trading_days,
116+
SUM(price_change) AS total_price_change,
117+
AVG(price_change) AS avg_daily_price_change,
118+
STDDEV(price_change) AS stddev_price_change,
119+
MIN(price_change) AS min_daily_change,
120+
MAX(price_change) AS max_daily_change,
121+
AVG(pct_change) AS avg_daily_pct_change,
122+
STDDEV(pct_change) AS stddev_pct_change,
123+
MIN(pct_change) AS min_daily_pct_change,
124+
MAX(pct_change) AS max_daily_pct_change,
125+
SUM(CASE WHEN price_change > 0 THEN 1 ELSE 0 END) AS positive_days,
126+
SUM(CASE WHEN price_change < 0 THEN 1 ELSE 0 END) AS negative_days,
127+
SUM(CASE WHEN price_change = 0 THEN 1 ELSE 0 END) AS neutral_days
128+
FROM filtered_data
129+
WHERE time_period != 'older'
130+
GROUP BY commodity_name, commodity_unit, time_period
131+
),
132+
133+
-- Combine aggregated results with period boundary prices
134+
combined_results AS (
135+
SELECT
136+
ar.*,
137+
sp.period_start_price,
138+
ep.period_end_price
139+
FROM aggregated_results AS ar
140+
LEFT JOIN start_prices AS sp
141+
ON ar.commodity_name = sp.commodity_name
142+
AND ar.time_period = sp.time_period
143+
LEFT JOIN end_prices AS ep
144+
ON ar.commodity_name = ep.commodity_name
145+
AND ar.time_period = ep.time_period
146+
),
147+
148+
-- Calculate final metrics
149+
final_metrics AS (
150+
SELECT
151+
*,
152+
CASE
153+
WHEN period_start_price > 0
154+
THEN (
155+
(period_end_price - period_start_price)
156+
/ period_start_price
157+
)
158+
* 100
159+
END AS total_period_return_pct,
160+
CASE
161+
WHEN trading_days > 0
162+
THEN (positive_days * 100.0) / trading_days
163+
END AS win_rate_pct,
164+
stddev_pct_change * SQRT(252) AS annualized_volatility_pct
165+
FROM combined_results
166+
)
167+
168+
-- Final results
169+
SELECT
170+
commodity_name,
171+
commodity_unit,
172+
time_period,
173+
period_start_date,
174+
period_end_date,
175+
trading_days,
176+
positive_days,
177+
negative_days,
178+
neutral_days,
179+
ROUND(total_period_return_pct, 2) AS total_return_pct,
180+
ROUND(avg_daily_pct_change, 4) AS avg_daily_return_pct,
181+
ROUND(annualized_volatility_pct, 2) AS volatility_pct,
182+
ROUND(win_rate_pct, 1) AS win_rate_pct,
183+
ROUND(total_price_change, 2) AS total_price_change,
184+
ROUND(avg_daily_price_change, 4) AS avg_daily_price_change,
185+
ROUND(min_daily_change, 2) AS worst_day_change,
186+
ROUND(max_daily_change, 2) AS best_day_change,
187+
ROUND(period_start_price, 2) AS period_start_price,
188+
ROUND(period_end_price, 2) AS period_end_price
189+
FROM final_metrics
190+
ORDER BY time_period, commodity_name
191+

0 commit comments

Comments
 (0)