Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Running heavy work: Airflow is not the compute engine

Airflow & DAGs · Operating Airflow in Production

Running heavy work: Airflow is not the compute engine

Easyairflow-60
orchestrationcompute-separationpandasarchitecture

Question

Should you run a 50 GB pandas transformation inside a PythonOperator?

Solution

Airflow is designed strictly as a workflow orchestrator, not a distributed compute engine. Running a 50 GB pandas transformation inside an Airflow PythonOperator violates this separation of concerns and frequently causes cluster outages.

What happens when workers execute heavy data tasks

Airflow worker processes run on shared virtual machines or Kubernetes pods provisioned for lightweight scheduling, API dispatch, and task coordination.

When a task attempts to load a 50 GB CSV or Parquet file into pandas:

  • Memory multiplication: Python and pandas objects carry significant memory overhead. A 50 GB dataset on disk often expands to 120 GB to 180 GB in memory during joins, grouping, and intermediate transformations.
  • The Out-of-Memory killer: The operating system's Linux kernel invokes the OOM (Out Of Memory) killer to protect host stability. It terminates the entire worker process with SIGKILL (exit code 137).
  • Collateral damage: When a Celery worker process is terminated, every other task instance running concurrently on that same worker is killed abruptly, corrupting pipeline states across unrelated DAGs.
  • Inefficient compute: Airflow worker nodes lack distributed memory pooling, query optimization, columnar partitioning, and data shuffle capabilities.

The right architectural pattern: orchestration versus compute

Airflow's responsibility is to coordinate workflows; dedicated data engines should execute the transformations:

Orchestrator (Airflow Worker):
  - Sends trigger payload to Spark cluster
  - Polls job status every 30 seconds
  - Logs URL to Spark UI
  - Receives success code and triggers downstream load
          ↓
Compute Engine (Databricks, EMR, BigQuery, Snowflake):
  - Distributes 50 GB dataset across 16 worker nodes
  - Executes parallel in-memory join and aggregation
  - Writes Parquet files directly to cloud storage

If the transformation is too small for a Spark cluster but still requires isolated compute, package the code into a container and launch it with KubernetesPodOperator. This isolates memory usage so that an OOM crash affects only that single pod without destabilizing the Airflow cluster.

PreviousNext