Hadoop clusters have always needed something to run jobs in the right order at the right time. For years that was Oozie, with workflows in XML and coordinators triggered by time and data. Most teams now use Apache Airflow instead, because pipelines are Python code, the UI shows every run, and the same scheduler can drive Hadoop, cloud services and databases. The difficulty is that Airflow and Hadoop are two systems with their own ideas about scheduling, security and failure, and most production incidents live in the gap between them.

This article explains how to put the pieces together: which Airflow components need to reach the cluster, how Kerberos tickets get to tasks, how to wait for data without wasting workers, how to submit Spark to YARN so that kills and retries behave, and how to make every daily partition safe to rerun. Examples use Airflow 3 imports and the current Apache Hadoop-related providers.

Advertisement

Two schedulers, one pipeline

Airflow decides when a task should run and records whether it succeeded. It does not run your heavy work itself. On Hadoop, the heavy work runs as YARN applications, and the ResourceManager decides when containers start based on queue capacity. So every Spark or Hive task has two lifecycles: the Airflow task instance, which is usually a client process such as spark-submit or a HiveServer2 session, and the YARN application that does the work.

Design around that fact. Airflow should know the YARN application's outcome, not just whether the client started. Killing an Airflow task should kill the YARN application. Retrying a task must not produce duplicate output. And Airflow's concurrency limits should match the capacity of the YARN queue, otherwise Airflow launches twenty jobs that then sit in YARN accepted state, each holding an Airflow worker slot and timing out.

Where each component runs

Airflow 3 splits into a scheduler, a DAG processor that parses DAG files, an API server that serves the UI and the task execution API, a triggerer for deferred tasks, workers, and a metadata database. Workers no longer talk to the metadata database directly; they get and report task state through the API server. That matters for Hadoop because workers on edge nodes inside a secured network need a route to the API server, not to the database.

Only the workers that run Hadoop tasks need Hadoop. Give them the Hadoop, Spark and Hive client binaries, the cluster's configuration directories (core-site.xml, hdfs-site.xml, yarn-site.xml, hive-site.xml), and a keytab. With the Celery executor, start those workers listening on a dedicated queue and set queue="hadoop" on Hadoop tasks; other tasks run on ordinary workers that never hold cluster credentials. Keep the scheduler and DAG processor free of Hadoop calls: code at the top level of a DAG file runs on every parse, so a WebHDFS listing there turns the DAG processor into a load generator against the NameNode.

Airflow orchestrates; YARN executes. Only edge-node workers carry Hadoop clients and keytabs.Schedulercreates task instancesDAG processorparses DAG filesAPI serverUI + task execution APIMetadata DBstate, historyQueue: default workerspure Python tasks, alertsQueue: hadoop (edge nodes)hadoop, spark, beeline clientsKerberos renewerairflow kerberos -> ccachetask APIticketWebHDFS / NameNodesensors: files and markersYARN ResourceManagerspark-submit, queue=etlHiveServer2 + metastorepartitions, checksSpark driver + executorscontainers on NodeManagerslaunchAirflow tracks the spark-submit process; YARN tracks the application. Both must agree on success, failure and kill.
The control plane never touches the cluster. Tasks marked for the hadoop queue run on edge-node workers that hold client configs and a renewed Kerberos ticket; YARN runs the actual work.
Advertisement

Kerberos: getting a ticket to every task

A secured cluster authenticates every HDFS, YARN and Hive call with Kerberos, explained in the Hadoop Kerberos article. Tasks run unattended for hours, and tickets expire, so Airflow ships a renewer: airflow kerberos reads a keytab and refreshes a credential cache on a fixed interval. Tasks find the ticket through the cache path, so the renewer and every client must agree on it.

# airflow.cfg on every edge-node worker
[kerberos]
principal = airflow/edge01.example.com@EXAMPLE.COM
keytab = /etc/security/keytabs/airflow.keytab
ccache = /run/airflow/krb5_ccache          # not /tmp: the default is often world-readable
reinit_frequency = 3600

# start the renewer next to the worker, and point clients at the same cache:
#   export KRB5CCNAME=/run/airflow/krb5_ccache
#   airflow kerberos &
#   airflow celery worker --queues hadoop

Run one renewer per edge-node worker host, or per container in Kubernetes. Recent Airflow versions also support airflow kerberos --one-time, which obtains a ticket once and exits, for a container that needs a ticket before its task starts. For Spark, prefer passing principal and keytab to the Spark operator so that long-running applications can obtain fresh delegation tokens on their own: a ticket in the client's cache does not help an executor that runs past the token lifetime. Restrict the keytab file to the Airflow user, and use a principal with access only to the paths and queues the pipelines need, enforced through Ranger or HDFS permissions.

The providers you actually use

Provider packageMain classesNotes
apache-airflow-providers-apache-hdfsWebHDFSHook, WebHdfsSensor, MultipleFilesWebHdfsSensorTalks to HDFS over WebHDFS. The old snakebite-based HDFSHook and HDFSSensor were removed in provider 4.0.0
apache-airflow-providers-apache-sparkSparkSubmitOperator, SparkSubmitHookWraps spark-submit; the connection's host sets the master (yarn) and its extras the deploy mode
apache-airflow-providers-apache-hiveHiveServer2Hook, HivePartitionSensor and friendsHiveServer2 for SQL; metastore-based sensors for partitions
standard operatorsBashOperatorLast resort for CLI tools such as distcp; you own exit codes and kill handling

WebHDFS runs over HTTP to the NameNode, so the sensor needs no Hadoop client libraries, only network access and SPNEGO authentication on secured clusters; the WebHDFS article covers the protocol. Check provider versions against your Airflow version when you upgrade: providers release independently, and the removal of the snakebite classes broke DAGs that had pinned nothing.

Worked example: a daily clickstream partition

An ingestion job writes one day of click logs to /data/landing/clicks/dt=YYYY-MM-DD and creates a _SUCCESS file when the day is complete. We want to sessionize the day with Spark, publish it as a Hive partition, and check that it has rows. The whole DAG fits on one screen.

from datetime import datetime, timedelta

from airflow.sdk import DAG, task
from airflow.providers.apache.hdfs.sensors.web_hdfs import WebHdfsSensor
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.providers.apache.hive.hooks.hive import HiveServer2Hook

LANDING = "/data/landing/clicks/dt={{ ds }}"
CURATED = "/data/curated/clicks"

with DAG(
    dag_id="clicks_daily",
    schedule="@daily",
    start_date=datetime(2026, 9, 1),
    catchup=False,
    max_active_runs=2,                       # bound backfill pressure on YARN
    default_args={
        "retries": 2,
        "retry_delay": timedelta(minutes=10),
        "queue": "hadoop",                   # run on edge-node workers only
        "pool": "yarn_etl",                  # cap concurrent jobs in the YARN queue
    },
) as dag:

    landed = WebHdfsSensor(
        task_id="wait_for_success_marker",
        filepath=LANDING + "/_SUCCESS",      # the producer writes this last
        webhdfs_conn_id="webhdfs_default",
        mode="reschedule",                   # free the worker slot between pokes
        pool="sensors",                      # sensors never take YARN job slots
        poke_interval=300,
        timeout=6 * 3600,
    )

    transform = SparkSubmitOperator(
        task_id="sessionize",
        application="hdfs:///apps/clicks/sessionize.py",
        conn_id="spark_yarn",                # host=yarn, deploy-mode=cluster in extras
        name="clicks_sessionize_{{ ds_nodash }}",
        conf={"spark.yarn.queue": "etl", "spark.dynamicAllocation.maxExecutors": "40"},
        application_args=["--input", LANDING,
                          "--staging", CURATED + "/_staging/dt={{ ds }}",
                          "--final", CURATED + "/dt={{ ds }}"],
        execution_timeout=timedelta(hours=2),
    )

    @task
    def register_partition(ds=None):
        hook = HiveServer2Hook(hiveserver2_conn_id="hiveserver2_default")
        hook.run(f"ALTER TABLE curated.clicks ADD IF NOT EXISTS "
                 f"PARTITION (dt='{ds}') LOCATION '{CURATED}/dt={ds}'")

    @task
    def check_rows(ds=None):
        hook = HiveServer2Hook(hiveserver2_conn_id="hiveserver2_default")
        rows = hook.get_records(f"SELECT count(*) FROM curated.clicks WHERE dt='{ds}'")
        if rows[0][0] == 0:
            raise ValueError(f"no rows for dt={ds}")   # fail loudly, retries will not help

    landed >> transform >> register_partition() >> check_rows()

Walk through the decisions. The sensor waits, in its own pool, for the marker file, not the directory, because a directory exists as soon as the first file arrives. It uses mode="reschedule", which releases the worker slot between checks; in the default poke mode, a sensor waiting six hours holds a worker for six hours, and a dozen such sensors can starve the pool. The pool yarn_etl has as many slots as the YARN queue can run jobs at once, so Airflow queues the excess instead of YARN. max_active_runs=2 stops a backfill of a month from submitting thirty jobs together.

The Spark job writes to a staging path and renames it into place, so a partially written attempt is never visible and a rerun replaces the day instead of appending to it.

# sessionize.py: write to staging, then swap into place, so a retry never leaves half a partition.
import argparse
from pyspark.sql import SparkSession

args = argparse.ArgumentParser()
for a in ("--input", "--staging", "--final"):
    args.add_argument(a, required=True)
opts = args.parse_args()

spark = SparkSession.builder.getOrCreate()
df = spark.read.parquet(opts.input)
sessions = df.groupBy("user_id", "session_id").agg({"ts": "min", "url": "count"})
sessions.coalesce(64).write.mode("overwrite").parquet(opts.staging)   # bounded file count

jvm = spark._jvm
fs = jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
final, staging = jvm.org.apache.hadoop.fs.Path(opts.final), jvm.org.apache.hadoop.fs.Path(opts.staging)
if fs.exists(final):
    fs.delete(final, True)          # rerun of the same day replaces the old output
if not fs.rename(staging, final):   # rename within one HDFS namespace is atomic
    raise RuntimeError("rename failed")

The partition registration uses ADD IF NOT EXISTS, so rerunning it is harmless. The row check raises a plain error, which fails the run visibly. Every task is keyed by the run's logical date ({{ ds }}), never by the current clock, so rerunning last Tuesday processes last Tuesday.

Making Spark on YARN behave inside Airflow

Use cluster deploy mode for production jobs. The driver then runs in a YARN container instead of on the edge node, so a busy or restarted Airflow worker does not take drivers with it, and memory on edge nodes stays predictable. In cluster mode, spark-submit by default keeps polling until the application finishes and exits with its status, which is what lets Airflow see the real outcome. Do not turn that wait off unless you add your own status check.

When an Airflow task is cleared, times out or is marked failed, the operator's kill handler runs. The Spark hook tracks the YARN application id from spark-submit output and tries to kill that application, but only if it captured the id and the kill command can authenticate. Verify this in a test: start a job, kill the task in the UI, and confirm in the ResourceManager that the application died. Orphaned applications are the classic symptom when it does not: the next retry starts while the old attempt is still writing the same output path.

Name applications with the DAG, task and date, as the example does. When an on-call engineer looks at YARN at 3 a.m., clicks_sessionize_20261001 is traceable and PySparkShell is not. Send each pipeline to a named YARN queue; the YARN queues article explains how capacity and preemption between queues work.

Moving from Oozie

Oozie conceptAirflow equivalentDifference to plan for
Workflow (XML DAG of actions)DAG in PythonLogic and loops become code; review it like code
Spark, Hive, shell actionsSparkSubmitOperator, HiveServer2Hook tasks, BashOperatorActions ran inside a YARN launcher; Airflow runs clients on its own workers
Coordinator with time frequencyschedule on the DAGLogical date replaces nominal time; check the timezone
Coordinator with input datasetsSensors, or asset-triggered DAGsAsset events are Airflow-internal; external data still needs a sensor
BundleNo direct equivalentGroup DAGs by tags and folders
SLA elementsexecution_timeout, alerts on failureDecide which delays page someone and which only warn

Migrate one pipeline end to end first, including Kerberos, kill handling and backfill, rather than translating every workflow mechanically. Run the old and new pipelines in parallel for a week, writing to different output paths, and compare row counts per partition before switching readers.

Failure modes

SymptomCauseFix
Tasks fail with GSS or no valid credentials after hoursRenewer not running on that worker, or ccache paths disagreeRun the renewer next to every Hadoop worker; set KRB5CCNAME consistently; alert on renewer failure
Worker slots exhausted, nothing runningSensors in poke mode waiting for late dataUse reschedule mode and timeouts; route sensors to a sized pool
Jobs stuck in YARN accepted state, then time outAirflow concurrency exceeds the YARN queuePools sized to the queue; max_active_runs; spread backfills
Duplicate or mixed data after retryJob appends or writes directly to the final pathStaging path plus atomic rename; overwrite by partition
Orphaned YARN applicationsKill handler could not find or kill the application idTest kills; cluster mode; name apps; periodic cleanup of stale applications
Slow DAG parsing, NameNode loadHDFS or Hive calls at DAG file top levelMove all I/O into tasks; keep DAG files declarative
Thousands of tiny files per daySpark writes one file per task partitionCoalesce or repartition before writing; see the small files article
Wrong day processed around midnightCode uses the current date or a different timezoneUse the logical date; set the DAG timezone explicitly

The small files point is easy to underestimate: a daily DAG that writes 2,000 files per partition creates 730,000 NameNode objects a year for one table. The small files article explains the cost and the compaction options.

Trade-offs

Airflow gives you code-defined pipelines, a strong UI, retries and backfills, and a single scheduler across Hadoop and everything else. It adds a second system to secure and operate, and it does not understand YARN capacity unless you encode it in pools. Oozie runs inside the cluster and shares its security model, which is simpler in an all-Hadoop shop, but it is hard to test and has little momentum. If your cluster is being retired in favour of cloud services, Airflow is the better bet because the same DAGs can switch from SparkSubmitOperator on YARN to managed Spark operators with small changes.

What to do next

  1. Inventory your pipelines and list, for each, its input marker, output path, YARN queue and the user it runs as.
  2. Stand up a dedicated hadoop queue of edge-node workers with client configs, a keytab and a running Kerberos renewer, and keep every other worker free of credentials.
  3. Create one pool per YARN queue with slots equal to the jobs that queue can run, and set max_active_runs on every DAG.
  4. Convert sensors to reschedule mode with explicit timeouts, waiting on success markers rather than directories.
  5. Make each Spark job write to staging and rename into place, and key every path by the logical date.
  6. Test kill and retry for one job end to end and confirm in the ResourceManager that no application survives the kill.
Key takeaway: Airflow schedules and YARN executes, so a reliable Hadoop pipeline makes the two agree: tasks run on edge-node workers that hold client configs and a renewed Kerberos ticket, pools match YARN queue capacity, sensors wait for success markers in reschedule mode, Spark runs in cluster mode with kills that reach YARN, and every job writes to staging and renames into a partition keyed by the logical date, so any day can be rerun safely.