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

text
+-------------+    +-------------+    +-------------+    +-------------+
|   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

python
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

python
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

python
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

python
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

python
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

python
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

python
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

python
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

python
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

python
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

python
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

python
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.


 

 
Â