How to forward Nomad logs and alert on failures¶
Enable forwarding on the configuration passed to an existing Nomad task
group:
cfg.forward_logs = True
cfg.log_chunk_size = 65536
For airflow-config YAML, set the same fields inside the task’s cfg mapping:
forward_logs: true
log_chunk_size: 65536
Grant read-fs capability in the job’s namespace to the Nomad CLI token on
Airflow workers. Keep Nomad task logging enabled. Open the check-job task
log to see stdout and stderr tagged with allocation ID, task name, and stream.
Logs are also drained before restart, stop, and purge; check those task logs
for final output. Each terminal drain reads at most sixteen chunks per stream.
Permission or allocation-file read errors produce warnings and leave job
health checks active.
Log chunks can split lines. A UTF-8 character split across chunk boundaries can appear as replacement characters in forwarded text; byte cursors still advance by the original byte count. Keep the default chunk size unless you need a smaller read bound.
Attach your existing alert callback when constructing the group:
from airflow_nomad import Nomad
Nomad(dag=dag, cfg=cfg, on_failure_callback=report_failed)
For YAML, use an importable function path on the task model:
on_failure_callback: alerts.report_failed
Failed lifecycle commands raise task exceptions. Job failure diagnostics
include allocation ID, task state, exit code, signal, and Nomad task events.
Airflow invokes failure callbacks when tasks exhaust retries. Configure
maxretrigger to bound workload recovery attempts; Nomad task restart and
reschedule policies apply independently.
How to monitor a persistent job between runs¶
Create a separate scheduled health DAG for a job that continues after its
management DAG finishes. Keep stop_on_exit=False on the management
configuration, then pass that configuration to check_nomad_health:
from airflow_pydantic import Dag, PythonTask
from airflow_nomad import check_nomad_health
watchdog = Dag(
dag_id="nomad-health",
schedule="*/5 * * * *",
start_date="2025-01-01",
catchup=False,
max_active_runs=1,
tasks={
"health": PythonTask(
python_callable=check_nomad_health,
op_kwargs={"cfg": cfg.model_dump()},
on_failure_callback=report_failed,
retries=0,
),
},
)
watchdog.instantiate()
For airflow-config, save config/nomad_health.yaml beside your DAG loader.
Supply the job configuration as op_kwargs.cfg:
# @package _global_
_target_: airflow_config.Configuration
_convert_: all
dags:
nomad-health:
schedule: "*/5 * * * *"
start_date: "2025-01-01"
catchup: false
max_active_runs: 1
tasks:
health:
_target_: airflow_pydantic.PythonTask
python_callable: airflow_nomad.check_nomad_health
on_failure_callback: alerts.report_failed
retries: 0
op_kwargs:
cfg:
forward_logs: true
job:
id: worker
namespace: default
type: service
task_groups:
- name: worker
tasks:
- name: worker
driver: exec
config:
command: /opt/jobs/worker
Save nomad_health.py in your DAG folder:
"""Generate Airflow DAGs for Nomad health checks."""
from airflow_config import load_config
load_config("config", "nomad_health").generate_in_mem()
The health callable checks the existing job without registering, restarting, or stopping
it. A missing, stopped, failed, or completed persistent job fails the health
task. Set op_kwargs.require_running: false if successfully completed batch
jobs should pass. Watchdog tasks keep byte cursors in nomad_log_offsets
XComs and read them from prior runs, including failed checks, to avoid replaying
retained output. Keep XCom history for this health task to preserve those
cursors. Limit the health DAG to one active run so checks do not overlap.
Direct calls without an Airflow task instance read a bounded snapshot
of retained logs on each invocation. A crash and
recovery entirely between checks may be missed; shorten the schedule to match
your detection requirement.