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
- Provider installieren und Airflow konfigurieren: Installation
- Architektur und Zuständigkeiten verstehen: Architektur
- Für einzelne Container-Tasks den MesosOperator verwenden.
- Für vollständige Beispiele DAG-Beispiele lesen.
- 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:
- Der Operator sendet den Container-Auftrag an
POST /v0/queue_command. - Der MesosExecutor reiht den Auftrag ein und nimmt ein passendes Mesos-Offer an.
- Der Operator fragt
GET /v0/task/<task_id>ab. - Der Operator wartet auf
TASK_FINISHEDoder 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
| Parameter | Beschreibung |
|---|---|
image | Container-Image; erforderlich. |
command | String oder Argumentliste. Strings laufen über /bin/sh -c. |
cpus | Angeforderte CPU-Ressourcen. |
mem_limit | Angeforderter Speicher, zum Beispiel 128m oder eine Zahl. memlimit bleibt als Alias verfügbar. |
disk | Angeforderter Mesos-Datenträger. |
environment | Dictionary mit Umgebungsvariablen. |
attributes | Liste von Mesos-Attributbedingungen. |
force_pull | Steuert, ob das Image erneut gezogen werden soll. |
network_mode | Docker-Netzwerkmodus. |
user | Benutzer im Container. |
volumes | Volume-Angaben. |
airflow_scheduler_url | URL der Executor-API; Standard ist operator_api_url beziehungsweise http://localhost:11000. |
poll_interval | Sekunden zwischen Statusabfragen. |
startup_timeout | Maximale 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:
- Läuft der Airflow-Scheduler mit
MesosExecutor? - Ist
airflow_scheduler_urlkorrekt und auf Port 11000 erreichbar? - Ist das Framework im Mesos-Master registriert?
- 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.