Migrating to duckdb
This commit is contained in:
@@ -0,0 +1,241 @@
|
||||
"""DuckDB-based silver and gold layer for the YNAB data pipeline."""
|
||||
|
||||
import os
|
||||
from datetime import date, timedelta
|
||||
|
||||
import duckdb
|
||||
|
||||
|
||||
class DuckDBLayer:
|
||||
"""Orchestrates silver (external base) and gold (views + materialized) layers via DuckDB."""
|
||||
|
||||
def __init__(self, base_path: str, warehouse_path: str, logger):
|
||||
self.base_path = base_path
|
||||
self.warehouse_path = warehouse_path
|
||||
self.logger = logger
|
||||
self.db = None
|
||||
os.makedirs(self.warehouse_path, exist_ok=True)
|
||||
self._init()
|
||||
|
||||
def _silver_parquet_expr(self, entity: str) -> str:
|
||||
return f"read_parquet('{self.base_path}/{entity}.parquet')"
|
||||
|
||||
def _init(self):
|
||||
self.db = duckdb.connect('/mnt/bulk/data/ynab_data.duckdb')
|
||||
self._setup_schemas()
|
||||
self._create_silver_external_tables()
|
||||
self._drop_gold()
|
||||
self._create_gold_views()
|
||||
self._create_gold_dates()
|
||||
self._create_gold_facts()
|
||||
|
||||
def _setup_schemas(self):
|
||||
self.db.execute("CREATE SCHEMA IF NOT EXISTS silver")
|
||||
self.db.execute("CREATE SCHEMA IF NOT EXISTS gold")
|
||||
|
||||
def _create_silver_external_tables(self):
|
||||
entities = ['accounts', 'categories', 'payees', 'transactions', 'scheduled_transactions']
|
||||
for entity in entities:
|
||||
sql = f"CREATE TABLE IF NOT EXISTS silver.{entity} AS SELECT * FROM {self._silver_parquet_expr(entity)}"
|
||||
try:
|
||||
self.db.execute(sql)
|
||||
self.logger.info(f"Created silver.{entity}")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create silver.{entity}: {e}")
|
||||
self.logger.info("Silver layer initialized")
|
||||
|
||||
def _drop_gold(self):
|
||||
for name in ['accounts', 'categories', 'payees']:
|
||||
self.db.execute(f"DROP VIEW IF EXISTS gold.{name}")
|
||||
self.db.execute("DROP TABLE IF EXISTS gold.dates")
|
||||
self.db.execute("DROP TABLE IF EXISTS gold.transactions")
|
||||
self.db.execute("DROP TABLE IF EXISTS gold.scheduled_transactions")
|
||||
|
||||
# ── gold views (dimensions) ─────
|
||||
|
||||
def _create_gold_accounts_view(self):
|
||||
self.logger.info("Transforming and creating gold.accounts view")
|
||||
sql = """
|
||||
SELECT
|
||||
id AS account_id,
|
||||
name AS account_name,
|
||||
"type" AS account_type,
|
||||
on_budget,
|
||||
closed,
|
||||
COALESCE(CAST(note AS VARCHAR), 'none') AS note,
|
||||
balance / 1000.0 AS balance,
|
||||
cleared_balance / 1000.0 AS cleared_balance,
|
||||
uncleared_balance / 1000.0 AS uncleared_balance,
|
||||
deleted
|
||||
FROM silver.accounts
|
||||
"""
|
||||
try:
|
||||
self.db.execute(f"CREATE VIEW gold.accounts AS {sql}")
|
||||
self.logger.info("Created gold.accounts view")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create gold.accounts view: {e}")
|
||||
|
||||
def _create_gold_categories_view(self):
|
||||
self.logger.info("Transforming and creating gold.categories view")
|
||||
sql = """
|
||||
SELECT
|
||||
id AS category_id,
|
||||
name AS category_name,
|
||||
category_group_name,
|
||||
hidden,
|
||||
COALESCE(CAST(note AS VARCHAR), 'none') AS note,
|
||||
budgeted / 1000.0 AS budgeted,
|
||||
activity / 1000.0 AS activity,
|
||||
balance / 1000.0 AS balance,
|
||||
deleted
|
||||
FROM silver.categories
|
||||
"""
|
||||
try:
|
||||
self.db.execute(f"CREATE VIEW gold.categories AS {sql}")
|
||||
self.logger.info("Created gold.categories view")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create gold.categories view: {e}")
|
||||
|
||||
def _create_gold_payees_view(self):
|
||||
self.logger.info("Transforming and creating gold.payees view")
|
||||
sql = """
|
||||
SELECT
|
||||
id AS payee_id,
|
||||
name AS payee_name,
|
||||
deleted
|
||||
FROM silver.payees
|
||||
"""
|
||||
try:
|
||||
self.db.execute(f"CREATE VIEW gold.payees AS {sql}")
|
||||
self.logger.info("Created gold.payees view")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create gold.payees view: {e}")
|
||||
|
||||
def _create_gold_views(self):
|
||||
self._create_gold_accounts_view()
|
||||
self._create_gold_categories_view()
|
||||
self._create_gold_payees_view()
|
||||
|
||||
# ── gold dates (materialized as Parquet, loaded into DuckDB) ─────
|
||||
|
||||
def _create_gold_dates(self):
|
||||
dates_path = os.path.join(self.warehouse_path, 'dates.parquet')
|
||||
self.logger.info("Creating gold.dates dimension")
|
||||
try:
|
||||
start = date(2020, 1, 1)
|
||||
end = date(2030, 12, 31)
|
||||
days = (end - start).days + 1
|
||||
rows = []
|
||||
for i in range(days):
|
||||
d = start + timedelta(days=i)
|
||||
month_z = str(d.month).zfill(2)
|
||||
day_z = str(d.day).zfill(2)
|
||||
rows.append({
|
||||
'date_id': f"{d.year}{month_z}{day_z}",
|
||||
'date': d,
|
||||
'year': int(d.year),
|
||||
'month': int(d.month),
|
||||
'day': int(d.day),
|
||||
'weekday': int(d.isoweekday()),
|
||||
'is_weekday': 1 if d.isoweekday() < 6 else 0,
|
||||
})
|
||||
|
||||
import pyarrow as pa
|
||||
import pyarrow.parquet as pq
|
||||
dates_df = pa.Table.from_pylist(rows)
|
||||
dates_df = dates_df.cast(pa.schema([
|
||||
('date_id', pa.string()),
|
||||
('date', pa.date32()),
|
||||
('year', pa.int32()),
|
||||
('month', pa.int8()),
|
||||
('day', pa.int8()),
|
||||
('weekday', pa.int8()),
|
||||
('is_weekday', pa.int8()),
|
||||
]))
|
||||
pq.write_table(dates_df, dates_path)
|
||||
|
||||
self.db.execute(f"CREATE TABLE gold.dates AS SELECT * FROM read_parquet('{dates_path}')")
|
||||
self.logger.info(f"Created gold.dates and persisted to {dates_path}")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create gold.dates: {e}")
|
||||
|
||||
# ── gold facts (materialized as Parquet) ─────
|
||||
|
||||
def _create_gold_transactions_fact(self):
|
||||
transactions_path = os.path.join(self.warehouse_path, 'transactions.parquet')
|
||||
self.logger.info("Transforming and persisting gold.transactions fact")
|
||||
sql = f"""
|
||||
COPY (
|
||||
SELECT
|
||||
memo,
|
||||
cleared,
|
||||
approved,
|
||||
COALESCE(CAST(flag_color AS VARCHAR), 'none') AS flag_color,
|
||||
account_id,
|
||||
payee_id,
|
||||
category_id,
|
||||
transfer_account_id,
|
||||
id AS transaction_id,
|
||||
CAST(YEAR(CAST(date AS DATE)) AS VARCHAR)
|
||||
|| LPAD(CAST(MONTH(CAST(date AS DATE)) AS VARCHAR), 2, '0')
|
||||
|| LPAD(CAST(DAY(CAST(date AS DATE)) AS VARCHAR), 2, '0') AS transaction_date,
|
||||
amount / 1000.0 AS transaction_amount,
|
||||
deleted
|
||||
FROM silver.transactions
|
||||
) TO '{transactions_path}' (FORMAT PARQUET)
|
||||
"""
|
||||
try:
|
||||
self.db.execute(sql)
|
||||
self.logger.info(f"Persisted transactions.fact to {transactions_path}")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create transactions.fact: {e}")
|
||||
|
||||
def _create_gold_scheduled_transactions_fact(self):
|
||||
st_path = os.path.join(self.warehouse_path, 'scheduled_transactions.parquet')
|
||||
self.logger.info("Transforming and persisting gold.scheduled_transactions fact")
|
||||
sql = f"""
|
||||
COPY (
|
||||
SELECT
|
||||
CAST(date_first AS DATE) AS date_first,
|
||||
CAST(date_next AS DATE) AS date_next,
|
||||
frequency,
|
||||
COALESCE(memo, 'none') AS memo,
|
||||
COALESCE(CAST(flag_color AS VARCHAR), 'none') AS flag_color,
|
||||
account_id,
|
||||
payee_id,
|
||||
category_id,
|
||||
transfer_account_id,
|
||||
id AS scheduled_transaction_id,
|
||||
amount / 1000.0 AS scheduled_transaction_amount
|
||||
FROM silver.scheduled_transactions
|
||||
) TO '{st_path}' (FORMAT PARQUET)
|
||||
"""
|
||||
try:
|
||||
self.db.execute(sql)
|
||||
self.logger.info(f"Persisted scheduled_transactions.fact to {st_path}")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to create scheduled_transactions.fact: {e}")
|
||||
|
||||
def _create_gold_facts(self):
|
||||
self._create_gold_transactions_fact()
|
||||
self._create_gold_scheduled_transactions_fact()
|
||||
# Create DuckDB tables from the persistent Parquet files so they're queryable
|
||||
self._register_facts()
|
||||
|
||||
def _register_facts(self):
|
||||
"""Register the fact Parquet files as DuckDB tables."""
|
||||
try:
|
||||
self.db.execute(f"CREATE TABLE gold.transactions AS SELECT * FROM read_parquet('{os.path.join(self.warehouse_path, 'transactions.parquet')}')")
|
||||
self.db.execute(f"CREATE TABLE gold.scheduled_transactions AS SELECT * FROM read_parquet('{os.path.join(self.warehouse_path, 'scheduled_transactions.parquet')}')")
|
||||
self.logger.info("Registered fact tables in DuckDB")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Failed to register fact tables: {e}")
|
||||
|
||||
def close(self):
|
||||
if self.db:
|
||||
self.db.close()
|
||||
|
||||
|
||||
def get_duckdb_layer(base_path: str, warehouse_path: str, logger) -> DuckDBLayer:
|
||||
"""Factory function to create a DuckDB layer and perform all transformations."""
|
||||
return DuckDBLayer(base_path, warehouse_path, logger)
|
||||
@@ -1,19 +1,19 @@
|
||||
'''Module to run the data pipeline'''
|
||||
|
||||
from pipeline import dimensions, facts, ingest, raw_to_base
|
||||
from pipeline import duckdb_layer, ingest, raw_to_base
|
||||
|
||||
|
||||
def pipeline_main(config, logger):
|
||||
'''Run the data pipeline'''
|
||||
logger.info('Starting data pipeline')
|
||||
|
||||
ingest.Ingest(config,logger).start_ingestion()
|
||||
raw_to_base.RawToBase(config,logger)
|
||||
dimensions.DimAccounts(config,logger)
|
||||
dimensions.DimCategories(config,logger)
|
||||
dimensions.DimPayees(config,logger)
|
||||
dimensions.DimDate(config,logger)
|
||||
facts.FactTransactions(config,logger)
|
||||
facts.FactScheduledTransactions(config,logger)
|
||||
ingest.Ingest(config, logger).start_ingestion()
|
||||
raw_to_base.RawToBase(config, logger)
|
||||
|
||||
duckdb_layer.get_duckdb_layer(
|
||||
config['base_data_path'],
|
||||
config['warehouse_data_path'],
|
||||
logger
|
||||
)
|
||||
|
||||
logger.info('Data pipeline completed successfully')
|
||||
|
||||
@@ -8,7 +8,6 @@ import polars as pl
|
||||
|
||||
import config.exit_codes as ec
|
||||
|
||||
#test comment for pr check
|
||||
|
||||
class RawToBase:
|
||||
def __init__(self, config: dict[str, Any],logger):
|
||||
|
||||
Reference in New Issue
Block a user