Apache Oozie is the workflow scheduler that most Hadoop clusters of the 2010s used to chain MapReduce, Hive, Sqoop and later Spark jobs into daily pipelines. It is also retired: the Apache Software Foundation retired Oozie in February 2025, and the last release was 5.2.1 in February 2021. You are rarely choosing Oozie today. You are keeping an Oozie estate alive, debugging a coordinator that stopped firing, or planning a migration and needing to know exactly what behaviour to reproduce.
This page explains the three definition layers, how the server, its database and the YARN launcher cooperate, how data triggers work, what the coordinator controls do to a backlog, and how retries and reruns behave, then works through a late-data incident and a checklist for running a scheduler that gets no more security fixes.
Three layers: workflow, coordinator, bundle
Oozie splits scheduling into three XML documents, each answering one question. A workflow answers what runs, in what order: a directed acyclic graph of action nodes (a Spark job, a Hive query, a shell script, an HDFS operation, an email) joined by control nodes. A workflow has no notion of time; you submit it and it runs once.
A coordinator answers when. It materializes one action per time slot (say, daily at 02:00 UTC), and each action starts the workflow only when its time has come and its input datasets exist. A bundle answers which coordinators belong together, starting, suspending and killing them as one unit. Job ids end in -W, -C or -B for the three layers, and a coordinator action id appends @N.
How a workflow actually executes
The Oozie server is a Java web application with a REST API (the CLI is a thin client over it, conventionally on port 11000). It runs no heavy work. When a workflow reaches an action node, the server starts a launcher: in Oozie 5 a YARN ApplicationMaster that pulls the action's jars from the shared library in HDFS and either runs the work itself (shell, java) or submits the real job (Spark, Hive). Older versions used a map-only MapReduce job as the launcher, which is why older clusters show a one-mapper job next to every real job.
The launcher calls back a unique URL on the server when the task completes; if that callback never arrives, the server falls back to polling. Every transition is written to the relational database behind the server (MySQL, PostgreSQL or Oracle in production; embedded Derby is only for experiments). That database is the scheduler: it holds every next materialization time, every waiting action and the history a rerun needs.
Workflow jobs move through PREP, RUNNING, SUSPENDED, SUCCEEDED, KILLED and FAILED. The example below is a typical daily workflow: a Spark ingest, then two independent tasks in parallel, with every error path routed to a notification and a <kill> node.
<workflow-app xmlns="uri:oozie:workflow:1.0" name="daily-sessions-wf">
<global>
<resource-manager>${resourceManager}</resource-manager>
<name-node>${nameNode}</name-node>
</global>
<start to="ingest"/>
<action name="ingest" retry-max="2" retry-interval="5">
<spark xmlns="uri:oozie:spark-action:1.0">
<prepare><delete path="${stagingDir}"/></prepare>
<master>yarn</master>
<mode>cluster</mode>
<name>sessions-ingest-${day}</name>
<class>com.example.SessionsIngest</class>
<jar>${appDir}/lib/sessions.jar</jar>
<spark-opts>--num-executors 20 --executor-memory 8G</spark-opts>
<arg>${inputDir}</arg>
<arg>${stagingDir}</arg>
</spark>
<ok to="split"/>
<error to="notify-failure"/>
</action>
<fork name="split">
<path start="hive-load"/>
<path start="metrics-export"/>
</fork>
<action name="hive-load"> ... <ok to="joined"/> <error to="notify-failure"/> </action>
<action name="metrics-export"> ... <ok to="joined"/> <error to="notify-failure"/> </action>
<join name="joined" to="end"/>
<action name="notify-failure">
<email xmlns="uri:oozie:email-action:0.2">
<to>data-oncall@example.com</to>
<subject>${wf:id()} failed at ${wf:lastErrorNode()}</subject>
<body>${wf:errorMessage(wf:lastErrorNode())}</body>
</email>
<ok to="fail"/>
<error to="fail"/>
</action>
<kill name="fail">
<message>Failed at [${wf:lastErrorNode()}]: ${wf:errorMessage(wf:lastErrorNode())}</message>
</kill>
<end name="end"/>
</workflow-app>Note the shape. <fork> and <join> must be used in pairs. Every action has exactly one <ok> and one <error> transition, so the only design question is where errors go. Workflow schema 1.0 uses <resource-manager> where older schemas used <job-tracker>, a common validation error when copying old examples. A <decision> node picks a transition from EL predicates, and EL gives failures context: ${wf:lastErrorNode()}, ${wf:errorMessage(node)}, and ${wf:run()} for the rerun count. Route every error to one notification that carries the job id, failed node and message, so on-call can go straight from the alert to oozie job -info.
Coordinators and data triggers
A coordinator's frequency is in minutes, but write it with EL: ${coord:days(1)} and ${coord:months(1)} understand calendar lengths, including daylight-saving days. Each nominal time produces one coordinator action.
The data side is what made Oozie popular. A dataset declares how often data appears and where, as a uri-template with ${YEAR}, ${MONTH}, ${DAY}, ${HOUR} placeholders. An input event names a range of dataset instances relative to the action's nominal time with ${coord:current(n)}. The action stays WAITING until every instance in the range is ready. Readiness is defined by the done-flag, and its three cases matter: if the element is omitted, Oozie waits for a _SUCCESS file in the directory; if it is present but empty, the directory's existence is enough; if it names a file, that file must exist.
<coordinator-app name="daily-sessions" frequency="${coord:days(1)}"
start="2026-01-01T02:00Z" end="2027-01-01T02:00Z" timezone="UTC"
xmlns="uri:oozie:coordinator:0.5">
<controls>
<timeout>360</timeout> <!-- minutes an action may wait for inputs -->
<concurrency>2</concurrency> <!-- at most 2 days running at once -->
<execution>FIFO</execution> <!-- oldest day first -->
<throttle>7</throttle> <!-- at most 7 actions WAITING -->
</controls>
<datasets>
<dataset name="clicks" frequency="${coord:hours(1)}"
initial-instance="2025-01-01T00:00Z" timezone="UTC">
<uri-template>${nameNode}/data/clicks/${YEAR}/${MONTH}/${DAY}/${HOUR}</uri-template>
<!-- omitted done-flag means: wait for _SUCCESS in the directory -->
</dataset>
<dataset name="sessions" frequency="${coord:days(1)}"
initial-instance="2025-01-01T00:00Z" timezone="UTC">
<uri-template>${nameNode}/data/sessions/${YEAR}/${MONTH}/${DAY}</uri-template>
</dataset>
</datasets>
<input-events>
<data-in name="input" dataset="clicks">
<start-instance>${coord:current(-24)}</start-instance>
<end-instance>${coord:current(-1)}</end-instance>
</data-in>
</input-events>
<output-events>
<data-out name="output" dataset="sessions">
<instance>${coord:current(-1)}</instance>
</data-out>
</output-events>
<action>
<workflow>
<app-path>${appDir}</app-path>
<configuration>
<property><name>inputDir</name><value>${coord:dataIn('input')}</value></property>
<property><name>day</name>
<value>${coord:formatTime(coord:nominalTime(), 'yyyy-MM-dd')}</value></property>
</configuration>
</workflow>
</action>
</coordinator-app>This coordinator runs daily at 02:00 UTC and needs the previous 24 hourly click partitions. ${coord:dataIn('input')} expands to the resolved directories, and coord:formatTime with coord:nominalTime() gives a date string. Prefer coord:current over coord:latest: current is fixed by the nominal time, while latest resolves to whatever is newest when evaluated, so a rerun can read different data.
The empty done-flag is a frequent cause of corrupt outputs: the writer creates the directory first and fills it over twenty minutes, and Oozie starts the moment it appears. Make producers write a marker last.
Concurrency, execution order and backlogs
The <controls> block decides what happens when many actions are due at once, which is exactly the situation after an outage or a late upstream feed. Coordinator actions move WAITING to READY (inputs present) to SUBMITTED to RUNNING and then to SUCCEEDED, FAILED or KILLED; a WAITING action whose timeout expires becomes TIMEDOUT, and some policies produce SKIPPED.
| Control | What it limits | Operational meaning |
|---|---|---|
| concurrency | actions RUNNING at once (default 1) | Raise it to catch up faster; each extra slot is another full workflow on the cluster. |
| execution | order among READY actions | FIFO (default) oldest first; LIFO newest first; LAST_ONLY skips a waiting or ready action once the next nominal time has passed; NONE skips one that is more than a tolerance past its nominal time. |
| throttle | actions in WAITING at once | Caps how far ahead materialization runs, which bounds the number of input checks the server performs. |
| timeout | minutes an action waits for inputs | After it, the action is TIMEDOUT and will not run unless rerun. Set it from the real lateness of the feed, not a guess. |
Choose execution order by what consumers need. A daily aggregate that reports depend on should be FIFO: each day must exist. A current-state snapshot is a LAST_ONLY job: after an outage, stale snapshots are worthless, and computing them delays today's run.
Timezones and daylight saving
Oozie processes coordinators in a fixed timezone without daylight saving, normally UTC. The documentation recommends region identifiers such as America/Los_Angeles over abbreviations, because PST does not shift to PDT. Run coordinators and datasets in UTC where you can; if a job must run at 06:00 local, dry-run the 23-hour and 25-hour days, where an hourly range such as current(-24) to current(-1) stops matching the business day.
Retries and reruns
Oozie has two very different recovery mechanisms. User-retry is declared on an action with retry-max and retry-interval (minutes), and since Oozie 4.3 retry-policy of periodic (the default) or exponential. The catch is in the specification: user-retry applies to a defined set of error codes, by default things such as an existing output directory, a missing file or an IOException in the action executor, and administrators extend the list with oozie.service.LiteWorkflowStoreService.user.retry.error.code.ext. Do not assume retry-max="2" retries every Spark failure; inject the failure you care about in a test environment and watch whether the action enters USER_RETRY.
Reruns are manual recovery. A workflow in an end state can be rerun under the same job id, either from the failed node (oozie.wf.rerun.failnodes=true) or by listing completed nodes to skip (oozie.wf.rerun.skip.nodes); only one option per rerun. Oozie does not clean up partial outputs for you, which is why every action that writes should start with a <prepare><delete/></prepare> of its own output. A coordinator rerun takes -action or -date ranges, deletes the output-event directories unless you pass -nocleanup, and re-reads the coordinator definition only with -refresh; -failed reruns just the failed workflow actions. Actions you deliberately abandon can be moved to IGNORED with oozie job -ignore so they stop showing as failures.
export OOZIE_URL=http://oozie.example.com:11000/oozie
oozie validate daily-sessions/coordinator.xml # schema check, no server state touched
oozie job -dryrun -config job.properties # materialize without running
oozie job -config job.properties -run # returns ...-C (coordinator id)
oozie job -info 0000123-260101000000000-oozie-oozi-C # actions and their states
# a workflow failed mid-way: rerun from the failed node, same workflow id
oozie job -config job.properties -rerun 0000456-260101000000000-oozie-oozi-W -D oozie.wf.rerun.failnodes=true
# a week of coordinator actions needs reprocessing with a fixed workflow.xml
oozie job -rerun 0000123-260101000000000-oozie-oozi-C -date 2026-09-20T02:00Z::2026-09-26T02:00Z -refresh
Operating the server
- Shared library. Action jars come from the sharelib in HDFS (
oozie.use.system.libpath=true). After upgrading Spark or Hive, install the new sharelib and runoozie admin -sharelibupdate; a version mismatch produces class errors that look like application bugs. - Database. Back it up, and keep the purge service enabled; an unpurged database makes every query and the web console progressively slower.
- High availability. Several servers can share one database behind a load balancer, coordinated through ZooKeeper, which makes ZooKeeper and the database the components whose outage stops every pipeline.
- Security. On a Kerberized cluster the server acts as a proxy user and Hive or HBase actions need credential blocks; see Kerberos on Hadoop.
- Logs. Server logs are searched per job with
-logfilter; the real job's logs are in YARN. How the YARN ApplicationMaster works explains why the launcher and the real job are separate applications.
Worked example: a late feed and a backlog
The click feed above stops for a day and a half because an upstream export broke. With timeout at 360 minutes, the 02:00 action for 22 September waits six hours and becomes TIMEDOUT; 23 September does the same. On the 24th the feed is repaired and the producer backfills 36 hours of partitions, each with its _SUCCESS marker.
TIMEDOUT is terminal, so the missed days do not recover on their own. The on-call engineer confirms with oozie job -info which actions timed out, checks that every hourly directory has its marker, and reruns the 22nd, waits for SUCCEEDED, then reruns the 23rd. The order matters because the 23rd reads the 22nd's sessions for carry-over visits, and with concurrency 2 a single two-day rerun would start both at once: FIFO orders submission, not completion. Afterwards the timeout was raised to the feed's observed lateness, an alert fires when an action is still WAITING at 04:00, and the previous day's sessions (current(-1) on the sessions dataset) became an input event so Oozie enforces the order itself.
Failure modes
- Silent waiting. An action waits forever because the done-flag never appears (the producer changed its marker name or path). Alert on actions WAITING past an expected time, not only on FAILED.
- Half-written inputs. An empty done-flag and a slow writer; the job succeeds on partial data. Always require a marker written last.
- Rerun reads different data.
coord:latestin input events, or a dataset whose upstream was rewritten. Usecoord:currentand immutable partitions. - Order mistaken for dependency. FIFO only orders submission; declare a day's dependency on the previous day as an input event.
- Launcher starvation. Launchers sharing a queue with real jobs can hold the capacity those jobs need. Give launchers a small separate queue.
- Unpatched dependencies. A retired project gets no CVE fixes; keep the server reachable only from trusted networks.
Leaving Oozie
Most teams move to Apache Airflow; Airflow on Hadoop covers the mapping, including Kerberos and Spark on YARN. Inventory first the semantics that do not translate one to one: done-flag triggers (sensors with explicit timeouts), the four execution orders (catchup and max active runs only approximate LAST_ONLY and LIFO), timezone materialization, and reruns with output cleanup. Run old and new schedulers in parallel against a shadow output path through at least one month-end. Sqoop actions need their own replacement (the Sqoop deep dive), and the cloud migration playbook places the scheduler in the wider move.
What to do next
- Inventory every bundle, coordinator and workflow: owner, frequency, timezone, controls, datasets and done-flags.
- Find every dataset with an empty done-flag and every input event using coord:latest, and fix them first.
- Add alerting on coordinator actions that stay WAITING past their expected start, not only on failures.
- Check that every writing action begins with a prepare delete of its own output so reruns are safe.
- Test user-retry with a deliberately injected failure to learn which errors your actions actually retry.
- Confirm the Oozie database is backed up, purged and restorable, and that the server is not exposed outside trusted networks.
- Start the migration inventory now, since the project receives no further releases.