🎉 75% of content is free forever — Unlock Premium from $10/mo →
CW
NEWSLIVESearch All Content
đŸ’ŧ Servicesâ„šī¸ Aboutâœ‰ī¸ ContactView Pricing Plansfrom $10

Performance Tuning and Optimization in Apache Airflow

đŸŸĸ Free Lesson

Advertisement

Performance Tuning and Optimization

Performance Optimization LayersDAG DesignTask granularitySchedulerHeartbeat tuningMetadata DBIndexing, poolingExecutorWorker scalingCachingDAG, Variable cacheScheduler Throughputmin(parse_rate, db_rate, dispatch_rate)DB OptimizationIndexes on dag_id, task_id, execution_dateWorker Scalingconcurrency * workers = max parallelismKey metrics: scheduler_heartbeat_interval, min_file_process_interval

Architecture Diagram

Formal Definitions

Detailed Explanation

Scheduler Optimization


Key Scheduler Settings:

ParameterDefaultRecommendedDescription
min_file_process_interval3030-60Seconds between DAG file scans
dag_dir_list_interval300300-600Seconds between directory listings
parsing_processes22-4Parallel DAG parsing processes
scheduler_heartbeat_sec55Scheduler heartbeat interval
parallelism3232-128Max concurrent tasks
max_active_tasks_per_dag1616-64Max tasks per DAG
max_active_runs_per_dag1616-32Max DAG runs per DAG
[scheduler]
min_file_process_interval = 30
parsing_processes = 2
parallelism = 32
max_active_tasks_per_dag = 16
store_serialized_dags = True

Database Optimization


Essential Indexes:

IndexTablePurpose
idx_task_instance_dag_runtask_instanceSpeeds up DAG run queries
idx_task_instance_statetask_instanceFast state filtering
idx_dag_run_statedag_runFast DAG run state queries

Connection Pool Settings:

ParameterRecommendedDescription
pool_size20Base connections
max_overflow30Extra connections for peaks
pool_timeout30Wait time for connection
pool_recycle1800Recycle after 30 min
pool_pre_pingTrueVerify connections
[database]
sql_alchemy_pool_size = 20
sql_alchemy_max_overflow = 30
sql_alchemy_pool_recycle = 1800
sql_alchemy_pool_pre_ping = True

Worker Optimization


Celery Worker Settings:

ParameterDescriptionRecommendation
worker_concurrencyTasks per worker4-16 (CPU-bound: 4-8, I/O-bound: 16-32)
worker_prefetch_multiplierTasks prefetched1 (fair scheduling)
worker_max_tasks_per_childRecycle worker after N tasks1000-2000
task_acks_lateAck after executionTrue (fault tolerance)
[celery]
worker_concurrency = 16
worker_prefetch_multiplier = 1
worker_max_tasks_per_child = 1000
task_acks_late = True

Tip: For CPU-bound tasks, use lower concurrency (4-8). For I/O-bound, use higher (16-32).

return pool_config

Architecture Diagram

### Worker Optimization

```python
# worker_optimization.py
import psutil
import os

def get_worker_recommendations():
    """Get resource recommendations based on system specs."""
    cpu_count = psutil.cpu_count()
    memory = psutil.virtual_memory()
    
    # Celery worker configuration
    worker_config = {
        # Concurrency = CPU cores (for CPU-bound tasks)
        # Concurrency = 2 * CPU cores (for I/O-bound tasks)
        'concurrency': min(cpu_count, 16),
        
        # Prefetch multiplier - how many tasks to prefetch
        'prefetch_multiplier': 1,
        
        # Maximum tasks per child before worker restart
        'max_tasks_per_child': 200,
        
        # Worker memory limit
        'max_memory_per_child': int(memory.total * 0.8 / cpu_count),
        
        # Task time limit (seconds)
        'task_time_limit': 3600,
        
        # Soft time limit (seconds) - raises SoftTimeLimitExceeded
        'task_soft_time_limit': 3000,
    }
    
    return worker_config

def monitor_worker_health():
    """Monitor worker health metrics."""
    import psutil
    
    metrics = {
        'cpu_percent': psutil.cpu_percent(interval=1),
        'memory_percent': psutil.virtual_memory().percent,
        'disk_usage': psutil.disk_usage('/').percent,
        'open_files': len(psutil.Process().open_files()),
        'connections': len(psutil.Process().connections()),
    }
    
    # Alert thresholds
    alerts = []
    if metrics['cpu_percent'] > 90:
        alerts.append(f"High CPU: {metrics['cpu_percent']}%")
    if metrics['memory_percent'] > 85:
        alerts.append(f"High Memory: {metrics['memory_percent']}%")
    if metrics['disk_usage'] > 90:
        alerts.append(f"High Disk: {metrics['disk_usage']}%")
    
    return {
        'metrics': metrics,
        'alerts': alerts,
        'healthy': len(alerts) == 0,
    }

Key Concepts Table

Optimization AreaMetricTargetImpact
DAG ParsingParse time< 1s per DAGHigh
Task LatencyQueue to start< 5sHigh
DB Query TimeAverage query< 100msHigh
Worker MemoryPer-worker< 4GBMedium
XCom SizePer operation< 48KBMedium
Log StorageDaily volume< 10GB/dayLow
Scheduler HeartbeatInterval5sLow

Code Examples

Performance Monitoring Dashboard

DAG Optimization Patterns

Resource-Aware Task Scheduling

Performance Metrics

Optimization Impact

OptimizationBeforeAfterImprovement
DAG Serialization10s parse2s parse80% faster
DB Indexing500ms query50ms query90% faster
Connection Pooling100ms connect10ms connect90% faster
Worker Concurrency4 tasks16 tasks4x throughput
XCom Backend500ms push50ms push90% faster

Resource Utilization

ResourceRecommendedWarningCritical
CPU< 70%70-85%> 85%
Memory< 70%70-85%> 85%
Disk I/O< 70%70-85%> 85%
Network< 50%50-80%> 80%
DB Connections< 70%70-85%> 85%

See Also

—
☆☆☆☆☆
0 ratings

Rate & Feedback

Need Expert Airflow Help?

Get personalized tutoring, project support, or professional consulting.

Advertisement