DAGs aus GitHub testen, synchronisieren und bereitstellen

Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)

In dieser Anleitung wird beschrieben, wie Sie eine CI/CD-Pipeline erstellen, um DAGs aus Ihrem GitHub-Repository in Ihrer Managed Airflow-Umgebung zu testen, zu synchronisieren und bereitzustellen.

Wenn Sie nur Daten aus anderen Diensten synchronisieren möchten, lesen Sie Daten aus anderen Diensten übertragen.

Übersicht über die CI/CD-Pipeline

Architekturdiagramm mit den Schritten des Ablaufs. Die Vorabüberprüfung und die PR-Überprüfung finden Sie im GitHub-Abschnitt und die DAG-Synchronisierung und die manuelle DAG-Überprüfung im Abschnitt „ Google Cloud “.
Abbildung 1. Architekturdiagramm mit den Schritten des Ablaufs (zum Vergrößern klicken)

Die CI/CD-Pipeline zum Testen, Synchronisieren und Bereitstellen von DAGs umfasst die folgenden Schritte:

  1. Sie nehmen eine Änderung an einem DAG vor und übertragen diese Änderung per Push an einen Entwicklungszweig in Ihrem Repository.

  2. Sie öffnen eine Pull-Anfrage für den Hauptzweig Ihres Repositorys.

  3. Cloud Build führt Einheitentests aus, um zu prüfen, ob Ihr DAG gültig ist.

  4. Ihre Pull-Anfrage wird genehmigt und mit dem Hauptzweig Ihres Repositorys zusammengeführt.

  5. Cloud Build synchronisiert Ihre Entwicklung Managed Airflow-Umgebung mit diesen neuen Änderungen.

  6. Sie prüfen, ob der DAG in Ihrer Entwicklung Umgebung wie erwartet funktioniert.

  7. Wenn Ihr DAG wie erwartet funktioniert, laden Sie ihn in Ihre Produktions- Managed Airflow-Umgebung hoch.

Ziele

  • Automatisierte Presubmit-Prüfung mit Cloud Build ausführen. Bei dieser Prüfung werden Einheitentests für einen DAG ausgeführt.
  • DAGs in Ihrer Managed Service for Apache Airflow-Entwicklungsumgebung mit DAGs in Ihrem GitHub-Repository synchronisieren.

Hinweis

  • In dieser Anleitung wird davon ausgegangen, dass Sie mit zwei identischen Managed Airflow-Umgebungen arbeiten: einer Entwicklungsumgebung und einer Produktionsumgebung.

    Für diese Anleitung konfigurieren Sie eine CI/CD-Pipeline nur für die Entwicklungsumgebung. Achten Sie darauf, dass die verwendete Umgebung keine Produktionsumgebung ist.

  • In dieser Anleitung wird davon ausgegangen, dass Sie Ihre DAGs und die zugehörigen Tests in einem GitHub-Repository gespeichert haben.

    Die Beispiel-CI/CD-Pipeline zeigt den Inhalt eines Beispiel-Repositorys. DAGs und Tests werden im Verzeichnis dags/ gespeichert. Die Dateien mit den Anforderungen, die Einschränkungsdatei und die Cloud Build-Konfigurationsdateien werden auf der obersten Ebene gespeichert. Das DAG-Synchronisierungsprogramm und seine Anforderungen befinden sich im Verzeichnis utils.

Presubmit-Prüfjob und Einheitentests erstellen

Der erste Cloud Build-Job führt eine Presubmit-Prüfung aus, bei der Einheitentests für Ihre DAGs ausgeführt werden.

Unittests hinzufügen

Erstellen Sie Einheitentests für Ihre DAGs, falls noch nicht geschehen. Speichern Sie diese Tests zusammen mit den DAGs in Ihrem Repository, jeweils mit dem Suffix _test. Die Testdatei für den DAG in example_dag.py ist beispielsweise example_dag_test.py. Diese Tests werden als Presubmit-Prüfung in Ihrem Repository ausgeführt.

Cloud Build-YAML-Konfiguration für die Presubmit-Prüfung erstellen

Erstellen Sie in Ihrem Repository eine YAML-Datei mit dem Namen test-dags.cloudbuild.yaml, die Ihren Cloud Build-Job für Presubmit-Prüfungen konfiguriert. Sie enthält drei Schritte:

  1. Installieren Sie die für Ihre DAGs erforderlichen Abhängigkeiten.
  2. Installieren Sie die für Ihre Einheitentests erforderlichen Abhängigkeiten.
  3. Führen Sie die DAG-Tests aus.

steps:
  # install dependencies
  - name: python:3.8-slim
    entrypoint: pip
    args: ["install", "-r", "requirements.txt", "-c", "constraints.txt", "--user"]

  - name: python:3.8-slim
    entrypoint: pip
    args: ["install", "-r", "requirements-test.txt", "--user"]

  # run in python 3.8 which is latest version in Cloud Composer
  - name: python:3.8-slim
    entrypoint: python3.8
    args: ["-m", "pytest", "-s", "dags/"]

Cloud Build-Trigger für die Presubmit-Prüfung erstellen

Folgen Sie der Anleitung zum Erstellen von Repositorys aus GitHub , um einen auf einer GitHub-App basierenden Trigger mit den folgenden Konfigurationen zu erstellen:

  • Name: test-dags

  • Ereignis: Pull-Anfrage

  • Quelle - Repository: Wählen Sie Ihr Repository aus.

  • Quelle - Basiszweig: ^main$ (ändern Sie main bei Bedarf in den Namen des Basiszweigs Ihres Repositorys)

  • Quelle - Kommentarsteuerung: nicht erforderlich

  • Build-Konfiguration - Cloud Build-Konfigurationsdatei: /test-dags.cloudbuild.yaml (der Pfad zu Ihrer Build-Datei)

DAG-Synchronisierungsjob erstellen und DAG-Dienstprogrammskript hinzufügen

Konfigurieren Sie als Nächstes einen Cloud Build-Job, der ein DAG-Dienstprogrammskript ausführt. Das Dienstprogrammskript in diesem Job synchronisiert Ihre DAGs mit Ihrer Managed Airflow-Umgebung, nachdem sie mit dem Hauptzweig in Ihrem Repository zusammengeführt wurden.

DAG-Dienstprogrammskript hinzufügen

Fügen Sie das DAG-Dienstprogrammskript Ihrem Repository hinzu. Dieses Dienstprogrammskript kopiert alle DAG-Dateien im Verzeichnis dags/ Ihres Repositorys in ein temporäres Verzeichnis und ignoriert alle Python-Dateien, die keine DAGs sind. Das Skript verwendet dann die Cloud Storage-Clientbibliothek, um alle Dateien aus diesem temporären Verzeichnis in das Verzeichnis dags/ im Bucket Ihrer Managed Airflow-Umgebung hochzuladen.

from __future__ import annotations

import argparse
import glob
import os
from shutil import copytree, ignore_patterns
import tempfile

# Imports the Google Cloud client library
from google.cloud import storage


def _create_dags_list(dags_directory: str) -> tuple[str, list[str]]:
    temp_dir = tempfile.mkdtemp()

    # ignore non-DAG Python files
    files_to_ignore = ignore_patterns("__init__.py", "*_test.py")

    # Copy everything but the ignored files to a temp directory
    copytree(dags_directory, f"{temp_dir}/", ignore=files_to_ignore, dirs_exist_ok=True)

    # The only Python files left in our temp directory are DAG files
    # so we can exclude all non Python files
    dags = glob.glob(f"{temp_dir}/*.py")
    return (temp_dir, dags)


def upload_dags_to_composer(
    dags_directory: str, bucket_name: str, name_replacement: str = "dags/"
) -> None:
    """
    Given a directory, this function moves all DAG files from that directory
    to a temporary directory, then uploads all contents of the temporary directory
    to a given cloud storage bucket
    Args:
        dags_directory (str): a fully qualified path to a directory that contains a "dags/" subdirectory
        bucket_name (str): the GCS bucket of the Cloud Composer environment to upload DAGs to
        name_replacement (str, optional): the name of the "dags/" subdirectory that will be used when constructing the temporary directory path name Defaults to "dags/".
    """
    temp_dir, dags = _create_dags_list(dags_directory)

    if len(dags) > 0:
        # Note - the GCS client library does not currently support batch requests on uploads
        # if you have a large number of files, consider using
        # the Python subprocess module to run gcloud storage cp --recursive on your dags
        # See https://cloud.google.com/storage/docs/gsutil/commands/cp for more info
        storage_client = storage.Client()
        bucket = storage_client.bucket(bucket_name)

        for dag in dags:
            # Remove path to temp dir
            dag = dag.replace(f"{temp_dir}/", name_replacement)

            try:
                # Upload to your bucket
                blob = bucket.blob(dag)
                blob.upload_from_filename(dag)
                print(f"File {dag} uploaded to {bucket_name}/{dag}.")
            except FileNotFoundError:
                current_directory = os.listdir()
                print(
                    f"{name_replacement} directory not found in {current_directory}, you may need to override the default value of name_replacement to point to a relative directory"
                )
                raise

    else:
        print("No DAGs to upload.")


if __name__ == "__main__":
    parser = argparse.ArgumentParser(
        description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
    )
    parser.add_argument(
        "--dags_directory",
        help="Relative path to the source directory containing your DAGs",
    )
    parser.add_argument(
        "--dags_bucket",
        help="Name of the DAGs bucket of your Composer environment without the gs:// prefix",
    )

    args = parser.parse_args()

    upload_dags_to_composer(args.dags_directory, args.dags_bucket)

Cloud Build-YAML-Konfiguration zum Synchronisieren von DAGs erstellen

Erstellen Sie in Ihrem Repository eine YAML-Datei mit dem Namen add-dags-to-composer.cloudbuild.yaml, die Ihren Cloud Build-Job zum Synchronisieren von DAGs konfiguriert. Sie enthält zwei Schritte:

  1. Installieren Sie die für das DAG-Dienstprogrammskript erforderlichen Abhängigkeiten.

  2. Führen Sie das Dienstprogrammskript aus, um die DAGs in Ihrem Repository mit Ihrer Managed Airflow-Umgebung zu synchronisieren.

steps:
  # install dependencies
  - name: python
    entrypoint: pip
    args: ["install", "-r", "utils/requirements.txt", "--user"]

  # run
  - name: python
    entrypoint: python
    args: ["utils/add_dags_to_composer.py", "--dags_directory=${_DAGS_DIRECTORY}", "--dags_bucket=${_DAGS_BUCKET}"]

Cloud Build-Trigger erstellen

Folgen Sie der Anleitung zum Erstellen von Repositorys aus GitHub , um einen auf einer GitHub-App basierenden Trigger mit den folgenden Konfigurationen zu erstellen:

  • Name: add-dags-to-composer

  • Ereignis: Per Push-Befehl an Zweig übertragen

  • Quelle - Repository: Wählen Sie Ihr Repository aus.

  • Quelle - Basiszweig: ^main$ (ändern Sie main bei Bedarf in den Namen des Basiszweigs Ihres Repositorys)

  • Quelle - Filter für enthaltene Dateien (glob): dags/**

  • Build-Konfiguration - Cloud Build-Konfigurationsdatei: /add-dags-to-composer.cloudbuild.yaml (der Pfad zu Ihrer Build-Datei)

Fügen Sie in der erweiterten Konfiguration zwei Substitutionsvariablen hinzu:

  • _DAGS_DIRECTORY - das Verzeichnis, in dem sich die DAGs in Ihrem Repository befinden. Wenn Sie das Beispiel-Repository aus dieser Anleitung verwenden, ist es dags/.

  • _DAGS_BUCKET – der Cloud Storage-Bucket, der das Verzeichnis dags/ in Ihrer Managed Airflow-Entwicklungsumgebung enthält. Lassen Sie das Präfix gs:// weg. Beispiel: us-central1-example-env-1234ab56-bucket.

CI/CD-Pipeline testen

In diesem Abschnitt folgen Sie einem DAG-Entwicklungsablauf, bei dem die neu erstellten Cloud Build-Trigger verwendet werden.

Presubmit-Job ausführen

Erstellen Sie eine Pull-Anfrage für Ihren Hauptzweig, um Ihren Build zu testen. Suchen Sie auf der Seite nach Ihrer Presubmit-Prüfung. Klicken Sie auf Details und wählen Sie Weitere Details in Google Cloud Build anzeigen aus, um Ihre Build-Logs in der Google Cloud Console aufzurufen.

eine lokale Airflow-Umgebung auszuführen.
Screenshot einer GitHub-Prüfung namens „test-dags“ mit einem roten Pfeil, der auf den Projektnamen in Klammern zeigt
Abbildung 2. Screenshot des Status der Cloud Build Presubmit-Prüfung auf GitHub (zum Vergrößern klicken)

Wenn Ihre Presubmit-Prüfung fehlgeschlagen ist, lesen Sie Build-Fehler beheben.

Prüfen, ob Ihr DAG in Ihrer Entwicklungsumgebung funktioniert

Nachdem Ihre Pull-Anfrage genehmigt wurde, führen Sie sie mit Ihrem Hauptzweig zusammen. Verwenden Sie die Google Cloud Console, um Ihre Build-Ergebnisse anzusehen. Wenn Sie viele Cloud Build-Trigger haben, können Sie Ihre Builds nach dem Triggernamen add-dags-to-composer filtern.

Nachdem der Cloud Build-Synchronisierungsjob erfolgreich ausgeführt wurde, wird der synchronisierte DAG in Ihrer Managed Airflow-Entwicklungsumgebung angezeigt. Dort können Sie prüfen, ob der DAG wie erwartet funktioniert.

DAG Ihrer Produktionsumgebung hinzufügen

Wenn der DAG wie erwartet funktioniert, fügen Sie ihn manuell Ihrer Produktionsumgebung hinzu. Laden Sie dazu die DAG-Datei in das Verzeichnis dags/ im Bucket Ihrer Managed Airflow Produktionsumgebung hoch.

Wenn Ihr DAG-Synchronisierungsjob fehlgeschlagen ist oder sich Ihr DAG in Ihrer Managed Airflow-Entwicklungsumgebung nicht wie erwartet verhält, lesen Sie Build-Fehler beheben.

Build-Fehler beheben

In diesem Abschnitt wird beschrieben, wie Sie häufige Szenarien für Build-Fehler beheben.

Was passiert, wenn meine Presubmit-Prüfung fehlgeschlagen ist?

Klicken Sie in Ihrer Pull-Anfrage auf Details und wählen Sie Weitere Details in Google Cloud Build anzeigen aus, um Ihre Build-Logs in der Google Cloud Console aufzurufen. Anhand dieser Logs können Sie das Problem mit Ihrem DAG beheben. Sobald Sie die Probleme behoben haben, übernehmen Sie die Korrektur und übertragen Sie sie per Push an Ihren Zweig. Die Presubmit-Prüfung wird noch einmal ausgeführt und Sie können die Logs als Debugging-Tool verwenden.

Was passiert, wenn mein DAG-Synchronisierungsjob fehlgeschlagen ist?

Verwenden Sie die Google Cloud Console, um Ihre Build-Ergebnisse anzusehen. Wenn Sie viele Cloud Build-Trigger haben, können Sie Ihre Builds nach dem Triggernamen add-dags-to-composer filtern. Prüfen Sie die Logs des Build-Jobs und beheben Sie die Fehler. Wenn Sie weitere Hilfe bei der Behebung der Fehler benötigen, nutzen Sie die Supportkanäle.

Was passiert, wenn mein DAG in meiner Managed Airflow-Umgebung nicht richtig funktioniert?

Wenn Ihr DAG in Ihrer Managed Airflow-Entwicklungsumgebung nicht wie erwartet funktioniert, übertragen Sie ihn nicht manuell in Ihre Managed Airflow-Produktionsumgebung. Führen Sie stattdessen einen der folgenden Schritte aus:

  • Machen Sie die Pull-Anfrage rückgängig mit den Änderungen, die Ihren DAG beschädigt haben, um ihn in den Zustand unmittelbar vor Ihren Änderungen zurückzuversetzen. Dadurch werden auch alle anderen Dateien in dieser Pull-Anfrage zurückgesetzt.
  • Erstellen Sie eine neue Pull-Anfrage, um die Änderungen am beschädigten DAG manuell zurückzusetzen.
  • Erstellen Sie eine neue Pull-Anfrage, um die Fehler in Ihrem DAG zu beheben.

Wenn Sie einen dieser Schritte ausführen, wird eine neue Presubmit-Prüfung und nach dem Zusammenführen der DAG-Synchronisierungsjob ausgelöst.

Nächste Schritte