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.
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.
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 hadoopRun 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 package | Main classes | Notes |
|---|---|---|
| apache-airflow-providers-apache-hdfs | WebHDFSHook, WebHdfsSensor, MultipleFilesWebHdfsSensor | Talks to HDFS over WebHDFS. The old snakebite-based HDFSHook and HDFSSensor were removed in provider 4.0.0 |
| apache-airflow-providers-apache-spark | SparkSubmitOperator, SparkSubmitHook | Wraps spark-submit; the connection's host sets the master (yarn) and its extras the deploy mode |
| apache-airflow-providers-apache-hive | HiveServer2Hook, HivePartitionSensor and friends | HiveServer2 for SQL; metastore-based sensors for partitions |
| standard operators | BashOperator | Last 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 concept | Airflow equivalent | Difference to plan for |
|---|---|---|
| Workflow (XML DAG of actions) | DAG in Python | Logic and loops become code; review it like code |
| Spark, Hive, shell actions | SparkSubmitOperator, HiveServer2Hook tasks, BashOperator | Actions ran inside a YARN launcher; Airflow runs clients on its own workers |
| Coordinator with time frequency | schedule on the DAG | Logical date replaces nominal time; check the timezone |
| Coordinator with input datasets | Sensors, or asset-triggered DAGs | Asset events are Airflow-internal; external data still needs a sensor |
| Bundle | No direct equivalent | Group DAGs by tags and folders |
| SLA elements | execution_timeout, alerts on failure | Decide 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
| Symptom | Cause | Fix |
|---|---|---|
| Tasks fail with GSS or no valid credentials after hours | Renewer not running on that worker, or ccache paths disagree | Run the renewer next to every Hadoop worker; set KRB5CCNAME consistently; alert on renewer failure |
| Worker slots exhausted, nothing running | Sensors in poke mode waiting for late data | Use reschedule mode and timeouts; route sensors to a sized pool |
| Jobs stuck in YARN accepted state, then time out | Airflow concurrency exceeds the YARN queue | Pools sized to the queue; max_active_runs; spread backfills |
| Duplicate or mixed data after retry | Job appends or writes directly to the final path | Staging path plus atomic rename; overwrite by partition |
| Orphaned YARN applications | Kill handler could not find or kill the application id | Test kills; cluster mode; name apps; periodic cleanup of stale applications |
| Slow DAG parsing, NameNode load | HDFS or Hive calls at DAG file top level | Move all I/O into tasks; keep DAG files declarative |
| Thousands of tiny files per day | Spark writes one file per task partition | Coalesce or repartition before writing; see the small files article |
| Wrong day processed around midnight | Code uses the current date or a different timezone | Use 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
- Inventory your pipelines and list, for each, its input marker, output path, YARN queue and the user it runs as.
- 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.
- Create one pool per YARN queue with slots equal to the jobs that queue can run, and set max_active_runs on every DAG.
- Convert sensors to reschedule mode with explicit timeouts, waiting on success markers rather than directories.
- Make each Spark job write to staging and rename into place, and key every path by the logical date.
- Test kill and retry for one job end to end and confirm in the ResourceManager that no application survives the kill.