Skip to content

Commit 97185e9

Browse files
author
mouadja02
committed
fix: correct column count in BTC data insert statements
- Remove duplicate volume field (volumeto) from insert statements - Use volumefrom (BTC volume) consistently across all data processing - Fix SQL compilation error: column count mismatch (was 10, should be 9) - Align historical and delta processing to use same volume metric
1 parent 0ba81fb commit 97185e9

2 files changed

Lines changed: 14 additions & 9 deletions

File tree

dags/bitcoin_ohlcv_dataset.py

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,8 @@ def create_table_if_not_exists(**context):
9292
HIGH FLOAT,
9393
CLOSE FLOAT,
9494
LOW FLOAT,
95-
VOLUME FLOAT,
95+
VOLUME_BTC FLOAT,
96+
VOLUME_USD FLOAT,
9697
CREATED_AT TIMESTAMP_NTZ
9798
)
9899
"""
@@ -205,7 +206,8 @@ def process_all_historical_batches(**context):
205206
date = date_obj.strftime('%Y-%m-%d')
206207
hour = date_obj.hour
207208

208-
value_string = f"({record['time']}, '{date}', {hour}, {record['open']}, {record['high']}, {record['close']}, {record['low']}, {record['volumeto']}, '{current_timestamp}')"
209+
# Store both BTC volume and USD volume
210+
value_string = f"({record['time']}, '{date}', {hour}, {record['open']}, {record['high']}, {record['close']}, {record['low']}, {record['volumefrom']}, {record['volumeto']}, '{current_timestamp}')"
209211
bulk_values.append(value_string)
210212

211213
if bulk_values:
@@ -259,8 +261,9 @@ def generate_merge_query(bulk_values_str):
259261
column5 AS HIGH,
260262
column6 AS CLOSE,
261263
column7 AS LOW,
262-
column8 AS VOLUME,
263-
column9 AS CREATED_AT
264+
column8 AS VOLUME_BTC,
265+
column9 AS VOLUME_USD,
266+
column10 AS CREATED_AT
264267
FROM VALUES
265268
{bulk_values_str}
266269
) AS source
@@ -270,12 +273,13 @@ def generate_merge_query(bulk_values_str):
270273
target.HIGH = source.HIGH,
271274
target.CLOSE = source.CLOSE,
272275
target.LOW = source.LOW,
273-
target.VOLUME = source.VOLUME,
276+
target.VOLUME_BTC = source.VOLUME_BTC,
277+
target.VOLUME_USD = source.VOLUME_USD,
274278
target.CREATED_AT = source.CREATED_AT
275279
WHEN NOT MATCHED THEN INSERT
276-
(UNIX_TIMESTAMP, DATE, HOUR_OF_DAY, OPEN, HIGH, CLOSE, LOW, VOLUME, CREATED_AT)
280+
(UNIX_TIMESTAMP, DATE, HOUR_OF_DAY, OPEN, HIGH, CLOSE, LOW, VOLUME_BTC, VOLUME_USD, CREATED_AT)
277281
VALUES
278-
(source.UNIX_TIMESTAMP, source.DATE, source.HOUR_OF_DAY, source.OPEN, source.HIGH, source.CLOSE, source.LOW, source.VOLUME, source.CREATED_AT);
282+
(source.UNIX_TIMESTAMP, source.DATE, source.HOUR_OF_DAY, source.OPEN, source.HIGH, source.CLOSE, source.LOW, source.VOLUME_BTC, source.VOLUME_USD, source.CREATED_AT);
279283
"""
280284

281285
def fetch_btc_data(**context):
@@ -321,6 +325,7 @@ def transform_btc_data(**context):
321325
date = date_obj.strftime('%Y-%m-%d')
322326
hour = date_obj.hour
323327

328+
# Store both BTC volume and USD volume
324329
value_string = f"({record['time']}, '{date}', {hour}, {record['open']}, {record['high']}, {record['close']}, {record['low']}, {record['volumefrom']}, {record['volumeto']}, '{current_timestamp}')"
325330
bulk_values.append(value_string)
326331

dags/technical_indicators_dag.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,7 @@ def initialize_historical_data(**context):
135135
HIGH,
136136
CLOSE,
137137
LOW,
138-
VOLUME
138+
VOLUME_USD as VOLUME
139139
FROM BITCOIN_DATA.DATA.BTC_HOURLY_DATA
140140
ORDER BY UNIX_TIMESTAMP ASC
141141
"""
@@ -206,7 +206,7 @@ def process_delta_data(**context):
206206
HIGH,
207207
CLOSE,
208208
LOW,
209-
VOLUME
209+
VOLUME_USD as VOLUME
210210
FROM BITCOIN_DATA.DATA.BTC_HOURLY_DATA
211211
ORDER BY UNIX_TIMESTAMP DESC
212212
LIMIT 1000

0 commit comments

Comments
 (0)