Setup DuckDB with Apache Airflow for automated data pipelines

Intermediate 45 min May 11, 2026 622 views
Ubuntu 24.04 Debian 12 AlmaLinux 9 Rocky Linux 9

Configure DuckDB as a high-performance analytical database backend for Apache Airflow workflows. Build automated data pipelines that process files, APIs, and databases using DuckDB's columnar engine.

Prerequisites

  • Root or sudo access
  • Python 3.8 or higher
  • 4GB RAM minimum
  • 10GB free disk space

What this solves

DuckDB provides a fast, embedded analytical database that works perfectly with Apache Airflow for data pipeline automation. This setup gives you columnar analytics on CSV files, Parquet data, and database connections without managing a separate database cluster. You'll configure the DuckDB provider for Airflow and build DAGs that automate data ingestion, transformation, and analysis.

Step-by-step installation

Update system packages

Start by updating your package manager to ensure you get the latest versions of Python and system libraries.

sudo apt update && sudo apt upgrade -y
sudo apt install -y python3 python3-pip python3-venv build-essential
sudo dnf update -y
sudo dnf install -y python3 python3-pip python3-devel gcc gcc-c++ make

Create dedicated user for Airflow

Run Airflow as a dedicated user for better security isolation and file permission management.

sudo useradd -m -s /bin/bash airflow
sudo usermod -aG sudo airflow
sudo su - airflow

Create Python virtual environment

Isolate Airflow and DuckDB dependencies from the system Python installation to avoid package conflicts.

python3 -m venv ~/airflow-venv
source ~/airflow-venv/bin/activate
pip install --upgrade pip setuptools wheel

Install Apache Airflow with DuckDB provider

Install Airflow with the DuckDB provider package and required dependencies for data pipeline operations.

export AIRFLOW_VERSION=2.8.1
export PYTHON_VERSION="$(python --version | cut -d" " -f 2 | cut -d"." -f 1-2)"
export CONSTRAINT_URL="https://raw.githubusercontent.com/apache/airflow/constraints-${AIRFLOW_VERSION}/constraints-${PYTHON_VERSION}.txt"
pip install "apache-airflow==${AIRFLOW_VERSION}" --constraint "${CONSTRAINT_URL}"
pip install apache-airflow-providers-duckdb
pip install duckdb pandas pyarrow

Initialize Airflow database

Set up the Airflow metadata database and create the default admin user for the web interface.

export AIRFLOW_HOME=~/airflow
airflow db init
airflow users create \
    --username admin \
    --firstname Admin \
    --lastname User \
    --role Admin \
    --email admin@example.com \
    --password SecurePassword123

Configure Airflow settings

Optimize Airflow configuration for DuckDB workflows and enable parallel task execution.

[core]
executor = LocalExecutor
max_active_runs_per_dag = 3
max_active_tasks_per_dag = 8
parallelism = 16

[scheduler]
dag_dir_list_interval = 300
max_threads = 2

[webserver]
expose_config = False
web_server_port = 8080
base_url = http://localhost:8080

Create DuckDB connection in Airflow

Configure a reusable DuckDB connection that your DAGs can reference for database operations.

mkdir -p ~/airflow/dags ~/airflow/logs ~/airflow/plugins ~/data
airflow connections add 'duckdb_default' \
    --conn-type 'duckdb' \
    --conn-host '/home/airflow/data/analytics.duckdb' \
    --conn-description 'Default DuckDB connection for analytics'

Create systemd service for Airflow scheduler

Set up Airflow scheduler to run automatically as a system service with proper logging and restart policies.

sudo tee /etc/systemd/system/airflow-scheduler.service > /dev/null << 'EOF'
[Unit]
Description=Airflow Scheduler
After=network.target
Wants=network.target

[Service]
User=airflow
Group=airflow
Type=simple
ExecStart=/home/airflow/airflow-venv/bin/airflow scheduler
Environment=AIRFLOW_HOME=/home/airflow/airflow
Restart=always
RestartSec=10
WorkingDirectory=/home/airflow
StandardOutput=journal
StandardError=journal

[Install]
WantedBy=multi-user.target
EOF

Create systemd service for Airflow webserver

Set up the Airflow web interface as a system service for DAG monitoring and management.

sudo tee /etc/systemd/system/airflow-webserver.service > /dev/null << 'EOF'
[Unit]
Description=Airflow Webserver
After=network.target
Wants=network.target

[Service]
User=airflow
Group=airflow
Type=simple
ExecStart=/home/airflow/airflow-venv/bin/airflow webserver --port 8080
Environment=AIRFLOW_HOME=/home/airflow/airflow
Restart=always
RestartSec=10
WorkingDirectory=/home/airflow
StandardOutput=journal
StandardError=journal

[Install]
WantedBy=multi-user.target
EOF

Enable and start Airflow services

Start both Airflow services and enable them to start automatically on system boot.

sudo systemctl daemon-reload
sudo systemctl enable --now airflow-scheduler airflow-webserver
sudo systemctl status airflow-scheduler airflow-webserver

Create your first DuckDB data pipeline

Create sample data processing DAG

Build a complete data pipeline that demonstrates DuckDB's capabilities with file processing, data transformation, and analytics.

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.duckdb.operators.duckdb import DuckDBOperator
from airflow.operators.python import PythonOperator
from airflow.providers.duckdb.hooks.duckdb import DuckDBHook
import pandas as pd
import os

default_args = {
    'owner': 'data-team',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'email_on_failure': True,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

dag = DAG(
    'duckdb_analytics_pipeline',
    default_args=default_args,
    description='DuckDB data analytics pipeline',
    schedule_interval=timedelta(hours=6),
    catchup=False,
    max_active_runs=1
)

def generate_sample_data(**context):
    """Generate sample sales data for processing"""
    import random
    from datetime import datetime, timedelta
    
    # Create sample data
    data = []
    base_date = datetime.now() - timedelta(days=30)
    
    for i in range(1000):
        data.append({
            'order_id': f'ORD-{i:06d}',
            'customer_id': f'CUST-{random.randint(1, 100):03d}',
            'product_id': f'PROD-{random.randint(1, 50):03d}',
            'quantity': random.randint(1, 10),
            'price': round(random.uniform(10, 500), 2),
            'order_date': (base_date + timedelta(days=random.randint(0, 30))).strftime('%Y-%m-%d'),
            'region': random.choice(['North', 'South', 'East', 'West'])
        })
    
    df = pd.DataFrame(data)
    os.makedirs('/home/airflow/data/raw', exist_ok=True)
    df.to_csv('/home/airflow/data/raw/sales_data.csv', index=False)
    print(f"Generated {len(data)} sales records")

generate_data = PythonOperator(
    task_id='generate_sample_data',
    python_callable=generate_sample_data,
    dag=dag
)

create_tables = DuckDBOperator(
    task_id='create_analytics_tables',
    duckdb_conn_id='duckdb_default',
    sql="""
        -- Create raw sales table
        CREATE TABLE IF NOT EXISTS raw_sales (
            order_id VARCHAR,
            customer_id VARCHAR,
            product_id VARCHAR,
            quantity INTEGER,
            price DECIMAL(10,2),
            order_date DATE,
            region VARCHAR
        );
        
        -- Create aggregated daily sales table
        CREATE TABLE IF NOT EXISTS daily_sales_summary (
            order_date DATE,
            region VARCHAR,
            total_orders INTEGER,
            total_revenue DECIMAL(12,2),
            avg_order_value DECIMAL(10,2),
            created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
        );
    """,
    dag=dag
)

load_data = DuckDBOperator(
    task_id='load_sales_data',
    duckdb_conn_id='duckdb_default',
    sql="""
        -- Clear existing data
        DELETE FROM raw_sales;
        
        -- Load CSV data into DuckDB
        INSERT INTO raw_sales 
        SELECT * FROM read_csv_auto('/home/airflow/data/raw/sales_data.csv');
        
        SELECT COUNT(*) as loaded_records FROM raw_sales;
    """,
    dag=dag
)

analyze_data = DuckDBOperator(
    task_id='analyze_sales_data',
    duckdb_conn_id='duckdb_default',
    sql="""
        -- Clear previous analysis
        DELETE FROM daily_sales_summary;
        
        -- Generate daily sales analytics
        INSERT INTO daily_sales_summary (order_date, region, total_orders, total_revenue, avg_order_value)
        SELECT 
            order_date,
            region,
            COUNT(*) as total_orders,
            ROUND(SUM(quantity * price), 2) as total_revenue,
            ROUND(AVG(quantity * price), 2) as avg_order_value
        FROM raw_sales
        GROUP BY order_date, region
        ORDER BY order_date DESC, region;
        
        -- Show top performing regions
        SELECT 
            region,
            SUM(total_orders) as total_orders,
            ROUND(SUM(total_revenue), 2) as total_revenue
        FROM daily_sales_summary
        GROUP BY region
        ORDER BY total_revenue DESC;
    """,
    dag=dag
)

def export_results(**context):
    """Export analysis results to files"""
    hook = DuckDBHook(duckdb_conn_id='duckdb_default')
    
    # Export daily summary
    daily_results = hook.get_pandas_df(
        "SELECT * FROM daily_sales_summary ORDER BY order_date DESC"
    )
    
    # Export regional summary
    regional_results = hook.get_pandas_df("""
        SELECT 
            region,
            SUM(total_orders) as total_orders,
            ROUND(SUM(total_revenue), 2) as total_revenue,
            ROUND(AVG(avg_order_value), 2) as avg_order_value
        FROM daily_sales_summary
        GROUP BY region
        ORDER BY total_revenue DESC
    """)
    
    os.makedirs('/home/airflow/data/reports', exist_ok=True)
    daily_results.to_csv('/home/airflow/data/reports/daily_sales.csv', index=False)
    regional_results.to_csv('/home/airflow/data/reports/regional_summary.csv', index=False)
    
    print(f"Exported {len(daily_results)} daily records and {len(regional_results)} regional summaries")

export_data = PythonOperator(
    task_id='export_analysis_results',
    python_callable=export_results,
    dag=dag
)

# Set task dependencies
generate_data >> create_tables >> load_data >> analyze_data >> export_data

Create advanced DuckDB analytics DAG

Build a more complex pipeline that demonstrates DuckDB's advanced features like JSON processing, window functions, and data export.

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.duckdb.operators.duckdb import DuckDBOperator
from airflow.operators.python import PythonOperator
from airflow.providers.duckdb.hooks.duckdb import DuckDBHook
import json
import os

default_args = {
    'owner': 'analytics-team',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'email_on_failure': False,
    'retries': 2,
    'retry_delay': timedelta(minutes=3)
}

dag = DAG(
    'duckdb_advanced_analytics',
    default_args=default_args,
    description='Advanced DuckDB analytics with JSON and time-series',
    schedule_interval=timedelta(hours=12),
    catchup=False
)

def create_json_data(**context):
    """Generate sample JSON event data"""
    import random
    from datetime import datetime, timedelta
    
    events = []
    base_time = datetime.now() - timedelta(hours=24)
    
    for i in range(500):
        event = {
            'event_id': f'evt_{i:06d}',
            'timestamp': (base_time + timedelta(minutes=random.randint(0, 1440))).isoformat(),
            'user_id': f'user_{random.randint(1, 100):03d}',
            'event_type': random.choice(['page_view', 'click', 'purchase', 'signup']),
            'properties': {
                'page': random.choice(['/home', '/products', '/cart', '/checkout']),
                'browser': random.choice(['chrome', 'firefox', 'safari', 'edge']),
                'device': random.choice(['desktop', 'mobile', 'tablet']),
                'value': round(random.uniform(0, 100), 2) if random.random() > 0.7 else None
            },
            'session_id': f'sess_{random.randint(1, 50):03d}'
        }
        events.append(event)
    
    os.makedirs('/home/airflow/data/events', exist_ok=True)
    with open('/home/airflow/data/events/user_events.json', 'w') as f:
        for event in events:
            f.write(json.dumps(event) + '\n')
    
    print(f"Generated {len(events)} event records")

generate_events = PythonOperator(
    task_id='generate_event_data',
    python_callable=create_json_data,
    dag=dag
)

create_event_tables = DuckDBOperator(
    task_id='create_event_tables',
    duckdb_conn_id='duckdb_default',
    sql="""
        -- Install and load JSON extension
        INSTALL json;
        LOAD json;
        
        -- Create events table
        CREATE TABLE IF NOT EXISTS user_events (
            event_id VARCHAR,
            timestamp TIMESTAMP,
            user_id VARCHAR,
            event_type VARCHAR,
            page VARCHAR,
            browser VARCHAR,
            device VARCHAR,
            value DECIMAL(10,2),
            session_id VARCHAR
        );
        
        -- Create hourly analytics table
        CREATE TABLE IF NOT EXISTS hourly_analytics (
            hour_bucket TIMESTAMP,
            event_type VARCHAR,
            device VARCHAR,
            event_count INTEGER,
            unique_users INTEGER,
            total_value DECIMAL(12,2),
            avg_value DECIMAL(10,2)
        );
        
        -- Create user session analytics
        CREATE TABLE IF NOT EXISTS session_analytics (
            session_id VARCHAR,
            user_id VARCHAR,
            session_start TIMESTAMP,
            session_end TIMESTAMP,
            session_duration_minutes INTEGER,
            page_views INTEGER,
            total_events INTEGER,
            conversion_value DECIMAL(10,2)
        );
    """,
    dag=dag
)

process_json_events = DuckDBOperator(
    task_id='process_json_events',
    duckdb_conn_id='duckdb_default',
    sql="""
        -- Clear existing event data
        DELETE FROM user_events;
        
        -- Load and parse JSON events
        INSERT INTO user_events
        SELECT 
            json_extract_string(json_data, '$.event_id') as event_id,
            STRPTIME(json_extract_string(json_data, '$.timestamp'), '%Y-%m-%dT%H:%M:%S.%f') as timestamp,
            json_extract_string(json_data, '$.user_id') as user_id,
            json_extract_string(json_data, '$.event_type') as event_type,
            json_extract_string(json_data, '$.properties.page') as page,
            json_extract_string(json_data, '$.properties.browser') as browser,
            json_extract_string(json_data, '$.properties.device') as device,
            CAST(json_extract_string(json_data, '$.properties.value') AS DECIMAL(10,2)) as value,
            json_extract_string(json_data, '$.session_id') as session_id
        FROM (
            SELECT json(line) as json_data 
            FROM read_text('/home/airflow/data/events/user_events.json')
        ) WHERE json_data IS NOT NULL;
        
        SELECT COUNT(*) as loaded_events FROM user_events;
    """,
    dag=dag
)

generate_hourly_analytics = DuckDBOperator(
    task_id='generate_hourly_analytics',
    duckdb_conn_id='duckdb_default',
    sql="""
        DELETE FROM hourly_analytics;
        
        INSERT INTO hourly_analytics
        SELECT 
            DATE_TRUNC('hour', timestamp) as hour_bucket,
            event_type,
            device,
            COUNT(*) as event_count,
            COUNT(DISTINCT user_id) as unique_users,
            COALESCE(SUM(value), 0) as total_value,
            ROUND(AVG(value), 2) as avg_value
        FROM user_events
        GROUP BY DATE_TRUNC('hour', timestamp), event_type, device
        ORDER BY hour_bucket DESC, event_count DESC;
        
        -- Show top analytics
        SELECT 
            hour_bucket,
            SUM(event_count) as total_events,
            SUM(unique_users) as total_unique_users,
            ROUND(SUM(total_value), 2) as total_revenue
        FROM hourly_analytics
        GROUP BY hour_bucket
        ORDER BY hour_bucket DESC
        LIMIT 10;
    """,
    dag=dag
)

analyze_user_sessions = DuckDBOperator(
    task_id='analyze_user_sessions',
    duckdb_conn_id='duckdb_default',
    sql="""
        DELETE FROM session_analytics;
        
        INSERT INTO session_analytics
        SELECT 
            session_id,
            user_id,
            MIN(timestamp) as session_start,
            MAX(timestamp) as session_end,
            CAST(EXTRACT(EPOCH FROM (MAX(timestamp) - MIN(timestamp)))/60 AS INTEGER) as session_duration_minutes,
            SUM(CASE WHEN event_type = 'page_view' THEN 1 ELSE 0 END) as page_views,
            COUNT(*) as total_events,
            COALESCE(SUM(CASE WHEN event_type = 'purchase' THEN value ELSE 0 END), 0) as conversion_value
        FROM user_events
        GROUP BY session_id, user_id
        HAVING COUNT(*) > 1
        ORDER BY conversion_value DESC, session_duration_minutes DESC;
        
        -- Show session insights
        SELECT 
            'Total Sessions' as metric,
            COUNT(*) as value
        FROM session_analytics
        UNION ALL
        SELECT 
            'Avg Session Duration (min)' as me

Automated install script

Run this to automate the entire setup

Vous ne voulez pas gérer cela vous-même ?

Nous gérons l'infrastructure des entreprises qui dépendent de leur disponibilité. Entièrement infogéré, avec un interlocuteur fixe qui connaît votre environnement.

Vous avez un interlocuteur fixe qui connaît votre installation

Rotterdam 05:58 · joignable par message, sans formulaire de ticket