1. LEARNING OBJECTIVES
By the end of this lesson, you will be able to:
-
Design complete ETL pipelines for financial data.
-
Implement data extraction from multiple sources (databases, APIs, files).
-
Apply data transformation logic for financial data.
-
Handle data quality and validation in ETL processes.
-
Implement incremental and full load strategies.
-
Build data integration workflows using Apache Airflow.
-
Monitor and troubleshoot ETL pipelines.
-
Design a real-time streaming ETL pipeline.
2. ETL ARCHITECTURE
2.1 ETL vs. ELT
| ETL (Extract-Transform-Load) | ELT (Extract-Load-Transform) |
|---|---|
| Transform before loading | Transform after loading |
| Optimized for data warehouses | Optimized for data lakes |
| Lower storage requirements | Higher storage requirements |
| More complex ETL processes | Leverages cloud data warehouse power |
2.2 ETL Pipeline Components
+-------------+ +-------------+ +-------------+ +-------------+
| Extract | -> | Validate | -> | Transform | -> | Load |
+-------------+ +-------------+ +-------------+ +-------------+
| | | |
v v v v
+-------------+ +-------------+ +-------------+ +-------------+
| Sources | | Quality | | Business | | Target |
| (DB, API, | | Checks | | Logic | | (DW, |
| Files) | | | | | | Data |
+-------------+ +-------------+ +-------------+ +-------------+
3. DATA EXTRACTION
3.1 Extracting from Relational Databases
import psycopg2 import pandas as pd from sqlalchemy import create_engine class DatabaseExtractor: """ Extracts data from relational databases. """ def __init__(self, connection_string): self.engine = create_engine(connection_string) def extract_table(self, table_name, columns=None, where=None, limit=None): """ Extract data from a table. """ query = f"SELECT {', '.join(columns) if columns else '*'} FROM {table_name}" if where: query += f" WHERE {where}" if limit: query += f" LIMIT {limit}" return pd.read_sql(query, self.engine) def extract_incremental(self, table_name, key_column, last_value, columns=None): """ Extract incremental data based on a key column. """ query = f""" SELECT {', '.join(columns) if columns else '*'} FROM {table_name} WHERE {key_column} > '{last_value}' ORDER BY {key_column} """ return pd.read_sql(query, self.engine) def extract_changes(self, table_name, last_etl_date): """ Extract changes using timestamp column. """ query = f""" SELECT * FROM {table_name} WHERE updated_at > '{last_etl_date}' OR created_at > '{last_etl_date}' """ return pd.read_sql(query, self.engine) # Usage extractor = DatabaseExtractor('postgresql://user:pass@localhost/source_db') new_orders = extractor.extract_incremental( 'orders', 'order_id', '2024-01-01' )
3.2 Extracting from APIs
import requests import json import time from typing import Dict, List, Optional class APIExtractor: """ Extracts data from REST APIs. """ def __init__(self, base_url, api_key, rate_limit=100): self.base_url = base_url self.api_key = api_key self.rate_limit = rate_limit self.request_count = 0 self.window_start = time.time() def _check_rate_limit(self): """Check and enforce rate limiting.""" now = time.time() if now - self.window_start > 60: self.request_count = 0 self.window_start = now self.request_count += 1 if self.request_count > self.rate_limit: wait_time = 60 - (now - self.window_start) if wait_time > 0: time.sleep(wait_time) self.request_count = 0 self.window_start = time.time() def get(self, endpoint, params=None, headers=None): """ Make a GET request with rate limiting. """ self._check_rate_limit() url = f"{self.base_url}{endpoint}" default_headers = { 'Authorization': f'Bearer {self.api_key}', 'Content-Type': 'application/json' } if headers: default_headers.update(headers) response = requests.get(url, params=params, headers=default_headers) response.raise_for_status() return response.json() def get_paginated(self, endpoint, params=None, page_param='page', limit_param='limit', limit=100): """ Extract paginated data. """ all_data = [] page = 1 while True: params = params or {} params[page_param] = page params[limit_param] = limit data = self.get(endpoint, params) if not data: break all_data.extend(data) page += 1 if len(data) < limit: break return all_data # Usage api_extractor = APIExtractor( base_url='https://api.example.com/v1', api_key='your_api_key', rate_limit=50 ) market_data = api_extractor.get_paginated( '/market_data', params={'symbol': 'AAPL', 'start_date': '2024-01-01'}, limit=100 )
3.3 Extracting from Files
import csv import json import pandas as pd import os from pathlib import Path class FileExtractor: """ Extracts data from various file formats. """ def __init__(self, base_path): self.base_path = Path(base_path) def extract_csv(self, file_path, delimiter=',', encoding='utf-8'): """ Extract data from CSV file. """ full_path = self.base_path / file_path return pd.read_csv(full_path, delimiter=delimiter, encoding=encoding) def extract_json(self, file_path): """ Extract data from JSON file. """ full_path = self.base_path / file_path with open(full_path, 'r') as f: return json.load(f) def extract_json_lines(self, file_path): """ Extract data from JSON Lines file (each line is a JSON object). """ full_path = self.base_path / file_path data = [] with open(full_path, 'r') as f: for line in f: data.append(json.loads(line.strip())) return pd.DataFrame(data) def extract_excel(self, file_path, sheet_name=0): """ Extract data from Excel file. """ full_path = self.base_path / file_path return pd.read_excel(full_path, sheet_name=sheet_name) def extract_parquet(self, file_path): """ Extract data from Parquet file. """ full_path = self.base_path / file_path return pd.read_parquet(full_path) def watch_directory(self, pattern, callback, interval=60): """ Watch a directory for new files. """ import time from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler class Handler(FileSystemEventHandler): def __init__(self, callback, pattern): self.callback = callback self.pattern = pattern def on_created(self, event): if not event.is_directory and event.src_path.endswith(self.pattern): self.callback(event.src_path) handler = Handler(callback, pattern) observer = Observer() observer.schedule(handler, str(self.base_path), recursive=False) observer.start() try: while True: time.sleep(interval) except KeyboardInterrupt: observer.stop() observer.join()
4. DATA TRANSFORMATION
4.1 Data Cleaning and Validation
import pandas as pd import numpy as np from datetime import datetime class DataTransformer: """ Applies transformations to financial data. """ def __init__(self): self.validation_rules = {} def add_validation_rule(self, column, rule_type, params): """ Add a validation rule for a column. """ if column not in self.validation_rules: self.validation_rules[column] = [] self.validation_rules[column].append({'type': rule_type, 'params': params}) def validate_data(self, df): """ Validate data against rules. """ errors = [] for column, rules in self.validation_rules.items(): for rule in rules: rule_type = rule['type'] params = rule['params'] if rule_type == 'not_null': null_count = df[column].isnull().sum() if null_count > 0: errors.append(f"{column}: {null_count} null values") elif rule_type == 'range': min_val = params.get('min') max_val = params.get('max') if min_val is not None: below_min = df[df[column] < min_val].shape[0] if below_min > 0: errors.append(f"{column}: {below_min} values below {min_val}") if max_val is not None: above_max = df[df[column] > max_val].shape[0] if above_max > 0: errors.append(f"{column}: {above_max} values above {max_val}") elif rule_type == 'in_set': allowed_values = params.get('values', []) invalid = df[~df[column].isin(allowed_values)].shape[0] if invalid > 0: errors.append(f"{column}: {invalid} invalid values") return errors def clean_numeric(self, df, column, fill_method='mean'): """ Clean numeric columns (handle missing values, outliers). """ if fill_method == 'mean': df[column] = df[column].fillna(df[column].mean()) elif fill_method == 'median': df[column] = df[column].fillna(df[column].median()) elif fill_method == 'zero': df[column] = df[column].fillna(0) return df def remove_outliers(self, df, column, method='iqr', threshold=1.5): """ Remove outliers from a column. """ if method == 'iqr': Q1 = df[column].quantile(0.25) Q3 = df[column].quantile(0.75) IQR = Q3 - Q1 lower_bound = Q1 - threshold * IQR upper_bound = Q3 + threshold * IQR df = df[(df[column] >= lower_bound) & (df[column] <= upper_bound)] elif method == 'zscore': z_scores = np.abs((df[column] - df[column].mean()) / df[column].std()) df = df[z_scores < threshold] return df def standardize_currencies(self, df, column, target_currency='USD', exchange_rates=None): """ Standardize currency values. """ if exchange_rates is None: exchange_rates = {'USD': 1.0, 'EUR': 1.08, 'GBP': 1.25, 'JPY': 0.0067} def convert(row): if row['currency'] == target_currency: return row[column] return row[column] * exchange_rates.get(row['currency'], 1.0) df[f'{column}_usd'] = df.apply(convert, axis=1) return df def calculate_derived_fields(self, df): """ Calculate common financial derived fields. """ # Trade value df['trade_value'] = df['quantity'] * df['price'] # Commission percentage if 'commission' in df.columns: df['commission_pct'] = df['commission'] / df['trade_value'] * 100 # P&L if 'sell_price' in df.columns and 'buy_price' in df.columns: df['pnl'] = (df['sell_price'] - df['buy_price']) * df['quantity'] # Date fields if 'trade_date' in df.columns: df['year'] = pd.to_datetime(df['trade_date']).dt.year df['month'] = pd.to_datetime(df['trade_date']).dt.month df['quarter'] = pd.to_datetime(df['trade_date']).dt.quarter df['day_of_week'] = pd.to_datetime(df['trade_date']).dt.dayofweek df['is_weekend'] = df['day_of_week'].isin([5, 6]) return df
4.2 Business Logic Transformations
class BusinessTransformer: """ Applies business-specific transformations. """ def apply_risk_rating(self, df): """ Apply risk rating to customers. """ def calculate_risk(row): score = 0 # Trading volume if row['total_volume'] > 1000000: score += 3 elif row['total_volume'] > 100000: score += 2 # Trade frequency if row['trade_count'] > 100: score += 2 elif row['trade_count'] > 20: score += 1 # Position size if row['max_position'] > 10000: score += 2 elif row['max_position'] > 1000: score += 1 # Risk tiers if score >= 7: return 'HIGH' elif score >= 4: return 'MEDIUM' else: return 'LOW' df['risk_rating'] = df.apply(calculate_risk, axis=1) return df def apply_portfolio_allocation(self, df): """ Calculate portfolio allocation metrics. """ # Total portfolio value per customer total_value = df.groupby('customer_id')['position_value'].sum() # Allocation percentage df['allocation_pct'] = df.apply( lambda row: (row['position_value'] / total_value.get(row['customer_id'], 1)) * 100, axis=1 ) return df def calculate_performance_metrics(self, df): """ Calculate performance metrics. """ # Sharpe ratio df['sharpe_ratio'] = (df['return_mean'] - df['risk_free_rate']) / df['return_std'] # Sortino ratio (downside deviation) downside_returns = df[df['returns'] < 0]['returns'] downside_std = downside_returns.std() df['sortino_ratio'] = (df['return_mean'] - df['risk_free_rate']) / downside_std # Information ratio df['information_ratio'] = (df['return_mean'] - df['benchmark_return']) / df['tracking_error'] return df
5. DATA LOADING
5.1 Loading to Data Warehouse
class DataLoader: """ Loads data into the data warehouse. """ def __init__(self, connection_string): self.engine = create_engine(connection_string) def load_full(self, df, table_name, if_exists='append'): """ Full load (truncate and reload). """ if if_exists == 'replace': df.to_sql(table_name, self.engine, if_exists='replace', index=False) else: df.to_sql(table_name, self.engine, if_exists='append', index=False) def load_incremental(self, df, table_name, key_column, batch_size=10000): """ Incremental load with upsert logic. """ conn = self.engine.raw_connection() cursor = conn.cursor() # Create temporary table temp_table = f"temp_{table_name}" cursor.execute(f"CREATE TEMP TABLE {temp_table} (LIKE {table_name})") # Insert data into temp table df.to_sql(temp_table, self.engine, if_exists='replace', index=False) # Merge using INSERT ... ON CONFLICT columns = ', '.join(df.columns) update_columns = ', '.join([f"{col} = EXCLUDED.{col}" for col in df.columns if col != key_column]) query = f""" INSERT INTO {table_name} ({columns}) SELECT {columns} FROM {temp_table} ON CONFLICT ({key_column}) DO UPDATE SET {update_columns} """ cursor.execute(query) conn.commit() cursor.close() conn.close() def load_with_history(self, df, table_name, key_column, history_column='valid_from'): """ Load with SCD Type 2 history tracking. """ conn = self.engine.raw_connection() cursor = conn.cursor() for _, row in df.iterrows(): # Check if record exists cursor.execute( f"SELECT {key_column} FROM {table_name} WHERE {key_column} = %s AND is_current = TRUE", (row[key_column],) ) exists = cursor.fetchone() if exists: # Close current record cursor.execute( f"UPDATE {table_name} SET valid_to = %s, is_current = FALSE " f"WHERE {key_column} = %s AND is_current = TRUE", (datetime.now().date(), row[key_column]) ) # Insert new record columns = ', '.join(df.columns) placeholders = ', '.join(['%s'] * len(df.columns)) query = f"INSERT INTO {table_name} ({columns}) VALUES ({placeholders})" cursor.execute(query, tuple(row.values)) conn.commit() cursor.close() conn.close()
5.2 Incremental Load Strategy
class IncrementalLoader: """ Manages incremental data loading. """ def __init__(self, connection_string, control_table='etl_control'): self.engine = create_engine(connection_string) self.control_table = control_table self._init_control_table() def _init_control_table(self): """Initialize the ETL control table.""" conn = self.engine.raw_connection() cursor = conn.cursor() cursor.execute(""" CREATE TABLE IF NOT EXISTS etl_control ( job_name VARCHAR(100) PRIMARY KEY, last_run TIMESTAMP, last_loaded_key VARCHAR(100), status VARCHAR(20), records_loaded INTEGER, error_message TEXT ) """) conn.commit() cursor.close() conn.close() def get_last_key(self, job_name): """Get the last loaded key for a job.""" conn = self.engine.raw_connection() cursor = conn.cursor() cursor.execute( f"SELECT last_loaded_key FROM {self.control_table} WHERE job_name = %s", (job_name,) ) result = cursor.fetchone() cursor.close() conn.close() return result[0] if result else None def update_control(self, job_name, last_key, status, records_loaded, error=None): """Update the ETL control table.""" conn = self.engine.raw_connection() cursor = conn.cursor() cursor.execute(""" INSERT INTO etl_control (job_name, last_run, last_loaded_key, status, records_loaded, error_message) VALUES (%s, %s, %s, %s, %s, %s) ON CONFLICT (job_name) DO UPDATE SET last_run = EXCLUDED.last_run, last_loaded_key = EXCLUDED.last_loaded_key, status = EXCLUDED.status, records_loaded = EXCLUDED.records_loaded, error_message = EXCLUDED.error_message """, (job_name, datetime.now(), last_key, status, records_loaded, error)) conn.commit() cursor.close() conn.close()
6. ETL ORCHESTRATION WITH APACHE AIRFLOW
6.1 Airflow DAG Definition
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.dummy import DummyOperator from datetime import datetime, timedelta import pandas as pd default_args = { 'owner': 'finance_team', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'email_on_failure': True, 'email_on_retry': False, 'retries': 3, 'retry_delay': timedelta(minutes=5) } def extract_trades(**context): """Extract trade data.""" from etl.extractors import DatabaseExtractor extractor = DatabaseExtractor('postgresql://user:pass@source/trading') last_run = context['last_run'] if 'last_run' in context else '2024-01-01' df = extractor.extract_incremental('trades', 'trade_id', last_run) # Push to XCom for downstream tasks context['task_instance'].xcom_push(key='trades_data', value=df.to_json()) return len(df) def transform_trades(**context): """Transform trade data.""" from etl.transformers import DataTransformer, BusinessTransformer # Pull data from XCom json_data = context['task_instance'].xcom_pull(key='trades_data') df = pd.read_json(json_data) transformer = DataTransformer() business = BusinessTransformer() # Apply transformations df = transformer.remove_outliers(df, 'amount') df = transformer.calculate_derived_fields(df) df = business.apply_risk_rating(df) context['task_instance'].xcom_push(key='transformed_trades', value=df.to_json()) return len(df) def load_trades(**context): """Load trade data.""" from etl.loaders import DataLoader json_data = context['task_instance'].xcom_pull(key='transformed_trades') df = pd.read_json(json_data) loader = DataLoader('postgresql://user:pass@warehouse/fintech') loader.load_incremental(df, 'fact_trade_execution', 'trade_id') return len(df) def validate_trades(**context): """Validate loaded data.""" # Validation logic pass # Define DAG dag = DAG( 'trade_etl_pipeline', default_args=default_args, description='Daily ETL pipeline for trades', schedule_interval='0 2 * * *', # Daily at 2 AM catchup=False ) # Define tasks start = DummyOperator(task_id='start', dag=dag) extract = PythonOperator( task_id='extract_trades', python_callable=extract_trades, dag=dag ) transform = PythonOperator( task_id='transform_trades', python_callable=transform_trades, dag=dag ) load = PythonOperator( task_id='load_trades', python_callable=load_trades, dag=dag ) validate = PythonOperator( task_id='validate_trades', python_callable=validate_trades, dag=dag ) end = DummyOperator(task_id='end', dag=dag) # Define dependencies start >> extract >> transform >> load >> validate >> end
6.2 Airflow Operators for Financial ETL
class FinancialETLOperator(BaseOperator): """ Custom operator for financial ETL tasks. """ template_fields = ('source_table', 'target_table', 'date_column') def __init__( self, source_table, target_table, date_column, connection_id='fintech_db', *args, **kwargs ): super().__init__(*args, **kwargs) self.source_table = source_table self.target_table = target_table self.date_column = date_column self.connection_id = connection_id def execute(self, context): from airflow.hooks.base import BaseHook conn = BaseHook.get_connection(self.connection_id) # Extract logic extract_sql = f""" SELECT * FROM {self.source_table} WHERE {self.date_column} >= '{{{{ prev_ds }}}}' AND {self.date_column} < '{{{{ ds }}}}' """ # Transform and load logic # ... self.log.info(f"ETL completed for {self.source_table} -> {self.target_table}")
7. REAL-TIME STREAMING ETL
7.1 Streaming ETL with Kafka and Spark
from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * class StreamingETL: """ Real-time streaming ETL using Spark Streaming. """ def __init__(self): self.spark = SparkSession.builder \ .appName("FinancialStreamingETL") \ .config("spark.sql.adaptive.enabled", "true") \ .getOrCreate() def create_stream(self, bootstrap_servers, topic): """ Create a streaming DataFrame from Kafka. """ return self.spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", bootstrap_servers) \ .option("subscribe", topic) \ .option("startingOffsets", "latest") \ .load() def parse_trade_data(self, df): """ Parse trade data from Kafka. """ schema = StructType([ StructField("trade_id", StringType(), True), StructField("symbol", StringType(), True), StructField("price", DoubleType(), True), StructField("quantity", IntegerType(), True), StructField("timestamp", TimestampType(), True) ]) return df.select( from_json(col("value").cast("string"), schema).alias("data") ).select("data.*") def transform_stream(self, df): """ Apply transformations to the stream. """ return df.withColumn("trade_value", col("price") * col("quantity")) \ .withColumn("year", year(col("timestamp"))) \ .withColumn("month", month(col("timestamp"))) \ .withColumn("day", dayofmonth(col("timestamp"))) def aggregate_stream(self, df): """ Aggregate streaming data. """ return df.groupBy( window(col("timestamp"), "1 minute"), col("symbol") ).agg( sum("quantity").alias("total_volume"), avg("price").alias("vwap"), count("*").alias("trade_count") ) def write_to_warehouse(self, df): """ Write streaming data to the warehouse. """ return df.writeStream \ .foreachBatch(self._batch_writer) \ .trigger(processingTime="10 seconds") \ .start() def _batch_writer(self, df, batch_id): """ Batch writer for streaming data. """ # Write to data warehouse df.write.jdbc( url="jdbc:postgresql://warehouse:5432/fintech", table="fact_streaming_trades", mode="append", properties={"user": "user", "password": "pass"} ) print(f"Batch {batch_id} written successfully") def run_pipeline(self, bootstrap_servers, topic): """ Run the complete streaming pipeline. """ stream = self.create_stream(bootstrap_servers, topic) parsed = self.parse_trade_data(stream) transformed = self.transform_stream(parsed) aggregated = self.aggregate_stream(transformed) query = self.write_to_warehouse(aggregated) return query
8. ETL MONITORING AND ERROR HANDLING
8.1 Error Handling Framework
import logging from datetime import datetime from functools import wraps class ETLMonitor: """ Monitors and logs ETL processes. """ def __init__(self, log_file='etl.log'): self.logger = logging.getLogger('ETLMonitor') self.logger.setLevel(logging.INFO) handler = logging.FileHandler(log_file) handler.setFormatter( logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') ) self.logger.addHandler(handler) def log_step(self, step_name): """ Decorator to log ETL steps. """ def decorator(func): @wraps(func) def wrapper(*args, **kwargs): self.logger.info(f"Starting {step_name}") try: start = datetime.now() result = func(*args, **kwargs) duration = (datetime.now() - start).total_seconds() self.logger.info(f"Completed {step_name} in {duration:.2f}s") return result except Exception as e: self.logger.error(f"Error in {step_name}: {str(e)}") raise return wrapper return decorator def send_alert(self, message, level='ERROR'): """ Send alert (email, Slack, etc.). """ self.logger.error(message) # Implementation for sending alerts pass # Usage monitor = ETLMonitor() @monitor.log_step('Extract') def extract_data(): # Extraction logic pass @monitor.log_step('Transform') def transform_data(): # Transformation logic pass
8.2 Data Quality Checks
class DataQualityChecker: """ Performs data quality checks on ETL outputs. """ def __init__(self): self.checks = [] def add_check(self, name, check_func, threshold): """ Add a data quality check. """ self.checks.append({ 'name': name, 'check': check_func, 'threshold': threshold }) def run_checks(self, df): """ Run all data quality checks. """ results = [] for check in self.checks: try: result = check['check'](df) passed = result <= check['threshold'] results.append({ 'name': check['name'], 'passed': passed, 'value': result, 'threshold': check['threshold'] }) except Exception as e: results.append({ 'name': check['name'], 'passed': False, 'error': str(e) }) return results def generate_report(self, results): """ Generate a data quality report. """ report = [] report.append("=" * 60) report.append("DATA QUALITY REPORT") report.append("=" * 60) for result in results: status = "PASS" if result['passed'] else "FAIL" report.append(f"{result['name']}: {status}") report.append(f" Value: {result.get('value', 'N/A')}") report.append(f" Threshold: {result.get('threshold', 'N/A')}") report.append("-" * 40) return "\n".join(report) # Example usage dq_checker = DataQualityChecker() def null_check(df): return df.isnull().sum().sum() / df.shape[0] * 100 def duplicate_check(df): return (df.duplicated().sum() / df.shape[0]) * 100 dq_checker.add_check('null_percentage', null_check, 5.0) dq_checker.add_check('duplicate_percentage', duplicate_check, 1.0) results = dq_checker.run_checks(df) print(dq_checker.generate_report(results))
9. SUMMARY FOR THE FINANCE PRACTITIONER
-
ETL Pipelines are the backbone of financial data warehousing.
-
Data Extraction must handle multiple sources (databases, APIs, files).
-
Data Transformation includes cleaning, validation, and business logic.
-
Data Loading strategies include full load, incremental load, and SCD.
-
Airflow provides orchestration for complex ETL workflows.
-
Streaming ETLÂ enables real-time financial analytics.
-
Monitoring and Error Handling are essential for production ETL systems.
Â