Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Airflow Mesos Provider

Der Airflow Mesos Provider integriert Apache Airflow mit Apache Mesos.

Er bietet zwei Ausführungsmodelle:

  • MesosExecutor: verteilt normale Airflow-Task-Workloads über ein Mesos-Framework.
  • MesosOperator: startet einen einzelnen DAG-Task als Mesos-Container und wartet auf dessen Abschluss.

Schnellstart

  1. Provider installieren und Airflow konfigurieren: Installation
  2. Architektur und Zuständigkeiten verstehen: Architektur
  3. Für einzelne Container-Tasks den MesosOperator verwenden.
  4. Für vollständige Beispiele DAG-Beispiele lesen.
  5. Mit Tests und Entwicklung verifizieren.

Voraussetzungen

  • Apache Airflow 3.x empfohlen; der Provider deklariert Airflow >=2.0.
  • Apache Mesos 1.6 oder neuer.
  • Python 3.x.
  • Für Docker-Container: ein Mesos-Agent mit aktivem Docker-Containerizer und Zugriff auf das gewünschte Image.

SSL und Mesos-Authentifizierung sind optional, werden für produktive Cluster aber empfohlen.

Weitere Dokumentation

Installation

Paketinstallation

pip install avmesos_airflow_provider

Der Provider benötigt außerdem eine kompatible Airflow-, avmesos- und HTTP-Umgebung. Für lokale Entwicklung stellt shell.nix eine reproduzierbare Umgebung bereit.

Airflow konfigurieren

Für die Ausführung normaler Airflow-Tasks über Mesos:

[core]
executor = avmesos_airflow_provider.executors.mesos_executor.MesosExecutor

Die vollständige Beispielkonfiguration steht in Konfiguration.

Lokale Entwicklung

nix-shell

Die Nix-Shell installiert Airflow, avmesos, den Provider und PostgreSQL. Sie richtet außerdem die lokale Airflow-Datenbank und die DAG-Umgebung ein.

MesosOperator verwenden

Der Operator benötigt keinen zweiten Executor. Der Airflow-Scheduler mit MesosExecutor stellt die interne API auf Port 11000 bereit. Ein Minimalbeispiel:

from datetime import datetime

from airflow import DAG
from avmesos_airflow_provider.operators.mesos import MesosOperator

with DAG("hello_mesos", schedule=None, start_date=datetime(2024, 1, 1), catchup=False) as dag:
    MesosOperator(
        task_id="hello",
        image="alpine:3.20",
        command="echo hello",
        cpus=0.1,
        mem_limit="128m",
    )

Details befinden sich in der Operator-Referenz.

Architektur

Der Provider stellt zwei unterschiedliche Ausführungswege bereit:

MesosExecutor

MesosExecutor ersetzt den Airflow-Executor. Airflow übergibt eingeplante Task-Workloads an ein Mesos-Framework. Das Framework nimmt Mesos-Offers an und startet daraus Airflow-Tasks als Container-Tasks.

Dieser Weg eignet sich, wenn ein gesamter Airflow-DAG oder viele normale Airflow-Tasks über Mesos verteilt werden sollen.

MesosOperator

MesosOperator läuft innerhalb eines normalen Airflow-DAGs und startet genau einen einzelnen Task als Mesos-Container. Er ist dem DockerOperator nachempfunden, verwendet aber die vom MesosExecutor bereitgestellte lokale API.

Der Ablauf ist:

  1. Der Operator sendet den Container-Auftrag an POST /v0/queue_command.
  2. Der MesosExecutor reiht den Auftrag ein und nimmt ein passendes Mesos-Offer an.
  3. Der Operator fragt GET /v0/task/<task_id> ab.
  4. Der Operator wartet auf TASK_FINISHED oder meldet einen terminalen Fehler an Airflow.

Die API läuft standardmäßig auf http://localhost:11000. Sie wird vom Scheduler-Prozess bereitgestellt und ist nicht die Mesos-Master-API auf Port 5050.

Datenfluss

Airflow Scheduler
      |
      | MesosExecutor API :11000
      v
MesosExecutor Framework
      |
      | Mesos scheduler protocol
      v
Mesos Master :5050
      |
      v
Mesos Agent -> Docker/Mesos Container

Der Operator erzeugt kein eigenes Mesos-Framework. Dadurch bleiben Offer-Verteilung, Ressourcenprüfung und Framework-Authentifizierung zentral im bestehenden Executor.

Konfiguration

Airflow-Executor

In airflow.cfg wird der Executor aktiviert:

[core]
executor = avmesos_airflow_provider.executors.mesos_executor.MesosExecutor

Mesos-Konfiguration

Die Werte sind Beispiele. Hosts, Zugangsdaten und Images müssen an die eigene Umgebung angepasst werden.

[mesos]
mesos_ssl = True
master = mesos-master.example.invalid:5050
framework_name = Airflow
checkpoint = True
failover_timeout = 604800
command_shell = True

task_cpu = 0.1
task_memory = 512
task_disk = 1000

authenticate = True
default_principal = <MESOS_PRINCIPAL>
default_secret = <MESOS_SECRET>

docker_image_slave = <AIRFLOW_RUNTIME_IMAGE>
docker_volume_driver = local
docker_volume_dag_name = airflowdags
docker_volume_dag_container_path = /airflow/dags/
docker_volume_logs_name = airflowlogs
docker_volume_logs_container_path = /airflow/logs/
docker_sock = /var/run/docker.sock
docker_user_group_id = <DOCKER_GROUP_ID>
docker_network_mode = bridge
docker_environment = []

api_username = <API_USERNAME>
api_password = <API_PASSWORD>
operator_api_url = http://localhost:11000

operator_api_url ist die Adresse, die der MesosOperator verwendet, wenn kein airflow_scheduler_url am Task gesetzt ist.

Attribute

Globale Attribute gelten für Executor-Tasks. Task-spezifische Attribute werden vom Operator beziehungsweise executor_config ergänzt:

mesos_attributes = ["airflow:true", "gpu:true?:cpu:true"]

Ein Operator kann zusätzlich Folgendes angeben:

MesosOperator(
    task_id="cpu_task",
    image="alpine:3.20",
    command="echo hello",
    attributes=["cpu:true"],
)

Keine Secrets oder produktiven Infrastrukturadressen in DAG-Dateien versionieren. Für Umgebungen sollten Airflow Connections, Variables oder externe Secret-Backends verwendet werden.

MesosOperator

MesosOperator führt einen einzelnen Container-Task unter Apache Mesos aus. Das Verhalten orientiert sich am Airflow DockerOperator, die Ausführung erfolgt jedoch über die MesosExecutor-API.

Beispiel

from datetime import datetime

from airflow import DAG
from avmesos_airflow_provider.operators.mesos import MesosOperator

with DAG(
    dag_id="mesos_operator_example",
    schedule=None,
    start_date=datetime(2024, 1, 1),
    catchup=False,
) as dag:
    MesosOperator(
        task_id="hello_mesos",
        image="alpine:3.20",
        command="echo hello from Mesos",
        cpus=0.1,
        mem_limit="128m",
        attributes=["airflow:true"],
    )

Parameter

ParameterBeschreibung
imageContainer-Image; erforderlich.
commandString oder Argumentliste. Strings laufen über /bin/sh -c.
cpusAngeforderte CPU-Ressourcen.
mem_limitAngeforderter Speicher, zum Beispiel 128m oder eine Zahl. memlimit bleibt als Alias verfügbar.
diskAngeforderter Mesos-Datenträger.
environmentDictionary mit Umgebungsvariablen.
attributesListe von Mesos-Attributbedingungen.
force_pullSteuert, ob das Image erneut gezogen werden soll.
network_modeDocker-Netzwerkmodus.
userBenutzer im Container.
volumesVolume-Angaben.
airflow_scheduler_urlURL der Executor-API; Standard ist operator_api_url beziehungsweise http://localhost:11000.
poll_intervalSekunden zwischen Statusabfragen.
startup_timeoutMaximale Wartezeit in Sekunden.

Airflow-Standardparameter wie task_id, retries, pool und queue werden über BaseOperator unterstützt.

Statusverhalten

Der Operator beendet sich erfolgreich bei:

TASK_FINISHED

Folgende Zustände führen zu AirflowException:

TASK_FAILED
TASK_ERROR
TASK_KILLED
TASK_LOST
TASK_UNREACHABLE

HTTP-Fehler, ungültige JSON-Antworten und das Überschreiten von startup_timeout werden ebenfalls als Task-Fehler gemeldet.

Einschränkung bei Abbruch

Die aktuelle Executor-API besitzt keinen separaten Kill-Endpunkt für direkt eingereihte Operator-Tasks. on_kill() protokolliert diese Einschränkung. Für lange Tasks sollten Airflow-Timeouts und kurze, kontrolliert abbrechbare Container-Kommandos verwendet werden.

DAG-Beispiele

MesosExecutor

Mit dem Executor werden normale Airflow-Tasks über Mesos verteilt. Die konkrete DAG-Task benötigt keine spezielle Operator-Klasse:

from datetime import datetime
from airflow import DAG
from airflow.providers.standard.operators.bash import BashOperator

with DAG("executor_example", schedule=None, start_date=datetime(2024, 1, 1), catchup=False) as dag:
    BashOperator(
        task_id="show_date",
        bash_command="date",
        executor_config={
            "cpus": 0.2,
            "mem_limit": "256m",
            "attributes": ["airflow:true"],
        },
    )

MesosOperator

Für genau einen Container-Task:

from datetime import datetime
from airflow import DAG
from avmesos_airflow_provider.operators.mesos import MesosOperator

with DAG("operator_example", schedule=None, start_date=datetime(2024, 1, 1), catchup=False) as dag:
    run = MesosOperator(
        task_id="run_command",
        image="alpine:3.20",
        command=["/bin/sh", "-c", "echo operator-ok && uname -a"],
        cpus=0.1,
        mem_limit="128m",
        environment={"EXAMPLE_MODE": "true"},
        attributes=["airflow:true"],
    )

Ein vollständiges, bewusst kurzes Test-DAG liegt unter docs/examples/dags/mesos_operator_test.py.

Entwicklung und Tests

Nix-Umgebung

Die Datei shell.nix stellt Python, Airflow, PostgreSQL und die Provider-Abhängigkeiten bereit:

nix-shell

Der Shell-Hook erstellt eine virtuelle Umgebung unter /tmp/python-dev, richtet die lokale Airflow-Datenbank ein und installiert den Provider editable.

Unit-Tests

Die Unit-Tests verwenden synthetische HTTP-Antworten und benötigen keinen Mesos- oder Airflow-Live-Dienst:

make test

Der Test-Target nutzt:

python3 -m unittest discover -s tests -v

Build

make build

Damit werden Source-Distribution und Wheel erzeugt. Vor dem Commit sollten mindestens make test, make build und git diff --check ausgeführt werden.

Live-Smoke-Test

Für eine echte Integration kann das Test-DAG in Airflow geladen und manuell gestartet werden. Voraussetzungen sind ein laufender Airflow-Scheduler mit MesosExecutor, eine erreichbare Executor-API auf Port 11000 und ein Mesos-Cluster mit passenden Agent-Attributen.

Die Mesos-Master-UI auf Port 5050 ist nur zur Beobachtung geeignet. Der Operator spricht nicht direkt mit der Master-UI.

Fehlerbehebung

Operator bleibt im Polling

Prüfen:

  1. Läuft der Airflow-Scheduler mit MesosExecutor?
  2. Ist airflow_scheduler_url korrekt und auf Port 11000 erreichbar?
  3. Ist das Framework im Mesos-Master registriert?
  4. Gibt es passende CPU-, Speicher- und Attribut-Offers?

Die Executor-API meldet HTTP 200 beim Queueing nur für die Annahme in die Warteschlange. Der Operator wartet danach weiter auf den Mesos-Status.

TASK_FAILED oder TASK_ERROR

Mesos-Agent-Logs und die Task-Details in der Mesos-UI prüfen. Häufige Ursachen sind ein nicht verfügbares Image, fehlende Ressourcen, nicht passende Attribute oder ein falscher Container-/Netzwerkmodus.

TASK_LOST

Der Agent oder die Framework-Verbindung ist verloren gegangen. Prüfen, ob das Framework neu verbunden ist und ob der Agent aktiv ist.

API antwortet mit 401

Die API-Konfiguration und der verwendete Endpunkt prüfen. /v0/dags ist geschützt; der Operator verwendet /v0/queue_command und /v0/task/<task_id>. Die API sollte nicht über einen öffentlichen Reverse Proxy ohne passende Authentifizierung veröffentlicht werden.

DAG wird nicht geladen

Zuerst den Import isoliert prüfen:

airflow dags list-import-errors

Danach sicherstellen, dass dags_folder auf das Verzeichnis mit der DAG-Datei zeigt und dass der Provider in derselben Python-Umgebung installiert ist wie Airflow.

Ressourcen passen nicht

cpus, mem_limit und disk müssen zu den freien Mesos-Offers passen. Bei Attributen muss mindestens ein aktiver Agent alle Bedingungen erfüllen. Globale Attribute aus mesos_attributes und task-spezifische Attribute werden zusammen verwendet.

Executor-API

Die API wird vom MesosExecutor im Airflow-Scheduler bereitgestellt. Standardadresse ist http://localhost:11000.

POST /v0/queue_command

Reiht einen direkten Container-Task ein.

Beispielkörper:

{
  "airflow_task_id": "airflow.example.hello",
  "container_type": "DOCKER",
  "command": ["/bin/sh", "-c", "echo hello"],
  "image": "alpine:3.20",
  "cpus": 0.1,
  "mem_limit": "128m",
  "attributes": ["airflow:true"],
  "environment": {"MODE": "test"}
}

Eine erfolgreiche Annahme liefert HTTP 200. Das bedeutet nur, dass der Auftrag in die Executor-Warteschlange übernommen wurde; die Mesos-Ausführung ist zu diesem Zeitpunkt noch nicht abgeschlossen.

GET /v0/task/<task_id>

Liefert den zuletzt bekannten Mesos-Status des Tasks. Der task_id muss URL-sicher kodiert werden, wenn er Zeichen außerhalb des üblichen Task-ID-Formats enthält.

Der Operator erwartet ein Objekt mit einem Statusfeld, zum Beispiel:

{
  "status": {
    "task_id": {"value": "airflow.example.hello"},
    "state": "TASK_FINISHED"
  }
}

Während der Ausführung können TASK_STAGING, TASK_STARTING und TASK_RUNNING auftreten. Terminale Fehlerzustände werden vom Operator in einen Airflow-Task-Fehler übersetzt.

Sicherheit

Die API sollte nur auf dem internen Airflow-/Scheduler-Netzwerk erreichbar sein. Zugangsdaten gehören in Airflow-Konfiguration oder Secret-Backends und nicht in versionierte DAGs. Die Mesos-Authentifizierung gegenüber dem Cluster wird separat über die [mesos]-Konfiguration gesteuert.