Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)
Auf dieser Seite wird beschrieben, wie Sie mit dem DataflowTemplateOperator Dataflow-Pipelines aus Managed Airflow starten.
Die Pipeline „Cloud Storage Text für BigQuery“ ist eine Batchpipeline, mit der Sie in Cloud Storage gespeicherte Textdateien hochladen, diese mit einer von Ihnen bereitgestellten benutzerdefinierten JavaScript-Funktion (User Defined Function, UDF) transformieren und die Ergebnisse in BigQuery ausgeben können.

Übersicht
Bevor Sie den Workflow starten, erstellen Sie die folgenden Einheiten:
Eine leere BigQuery-Tabelle aus einem leeren Dataset, die die folgenden Spalten mit Informationen enthält:
location,average_temperature,monthund optionalinches_of_rain,is_currentundlatest_measurement.Eine JSON-Datei, mit der die Daten aus der Datei
.txtin das richtige Format für das Schema der BigQuery-Tabelle normalisiert werden. Das JSON-Objekt enthält ein Array vonBigQuery Schema. Jedes Objekt enthält einen Spaltennamen, einen Eingabetyp und gibt an, ob es sich um ein Pflichtfeld handelt.Eine
.txt-Eingabedatei, die die Daten enthält, die per Batch in die BigQuery-Tabelle hochgeladen werden.Eine in JavaScript geschriebene benutzerdefinierte Funktion, die jede Zeile der
.txt-Datei in die für unsere Tabelle relevanten Variablen umwandelt.Eine Airflow-DAG-Datei, die auf den Speicherort dieser Dateien verweist.
Als Nächstes laden Sie die Datei
.txt, die UDF-Datei.jsund die Schemadatei.jsonin einen Cloud Storage-Bucket hoch. Sie laden die DAG auch in Ihre Managed Airflow-Umgebung hoch.Nachdem der DAG hochgeladen wurde, führt Airflow einen Task daraus aus. Bei dieser Aufgabe wird eine Dataflow-Pipeline gestartet, in der die benutzerdefinierte Funktion auf die Datei
.txtangewendet und gemäß dem JSON-Schema formatiert wird.Schließlich werden die Daten in die BigQuery-Tabelle hochgeladen, die Sie zuvor erstellt haben.
Hinweis
- Für diesen Leitfaden sind JavaScript-Kenntnisse erforderlich, um die benutzerdefinierte Funktion zu schreiben.
- In diesem Leitfaden wird davon ausgegangen, dass Sie bereits eine Managed Airflow-Umgebung haben. Informationen zum Erstellen einer Umgebung finden Sie unter Umgebung erstellen. Sie können mit dieser Anleitung eine beliebige Version von Managed Airflow verwenden.
Aktivieren Sie die APIs für Managed Service for Apache Airflow, Dataflow, Cloud Storage und BigQuery, falls sie noch nicht aktiviert sind.
Rollen, die zum Aktivieren von APIs erforderlich sind
Zum Aktivieren von APIs benötigen Sie die Berechtigung
serviceusage.services.enable. Wenn Sie das Projekt erstellt haben, haben Sie diese Berechtigung wahrscheinlich bereits über die Rolle „Inhaber“ (roles/owner). Andernfalls können Sie diese Berechtigung über die Rolle „Service Usage-Administrator“ (roles/serviceusage.serviceUsageAdmin) erhalten. Informationen zum Zuweisen von RollenPrüfen Sie, ob Sie die folgenden Berechtigungen haben:
- Managed Airflow-Rollen: Sie können eine Umgebung erstellen (falls Sie noch keine haben), Objekte im Bucket der Umgebung verwalten, DAGs ausführen und auf die Airflow-UI zugreifen.
- Cloud Storage-Rollen: Erstellen Sie einen Bucket und verwalten Sie Objekte darin.
- BigQuery-Rollen: Dataset und Tabelle erstellen, Daten in der Tabelle ändern, Tabellenschema und ‑metadaten ändern.
- Dataflow-Rollen: Dataflow-Jobs ansehen.
Achten Sie darauf, dass das Dienstkonto Ihrer Umgebung die Berechtigungen zum Erstellen von Dataflow-Jobs, zum Zugriff auf den Cloud Storage-Bucket und zum Lesen und Aktualisieren von Daten für die Tabelle in BigQuery hat.
Leere BigQuery-Tabelle mit einer Schemadefinition erstellen
Erstellen Sie eine BigQuery-Tabelle mit einer Schemadefinition. Sie verwenden diese Schemadefinition später in diesem Leitfaden. In dieser BigQuery-Tabelle werden die Ergebnisse des Batch-Uploads gespeichert.
So erstellen Sie eine leere Tabelle mit einer Schemadefinition:
Console
Rufen Sie in der Google Cloud Console die Seite „BigQuery“ auf:
Maximieren Sie im Navigationsbereich im Abschnitt Ressourcen Ihr Projekt.
Klicken Sie im Detailbereich auf Dataset erstellen.

Geben Sie auf der Seite „Dataset erstellen“ im Abschnitt Dataset-ID den Namen
average_weatherfür Ihr Dataset ein. Übernehmen Sie für alle anderen Felder die Standardeinstellungen.
Klicken Sie auf Dataset erstellen.
Kehren Sie zum Navigationsbereich zurück und maximieren Sie im Abschnitt Ressourcen Ihr Projekt. Klicken Sie dann auf das Dataset
average_weather.Klicken Sie im Detailfeld auf Tabelle erstellen.

Wählen Sie auf der Seite Tabelle erstellen im Abschnitt Quelle die Option Leere Tabelle aus.
Gehen Sie auf der Seite Create table (Tabelle erstellen) im Abschnitt Destination (Ziel) so vor:
Wählen Sie für Dataset-Name das Dataset
average_weatheraus.
Geben Sie im Feld Tabellenname den Namen
average_weatherein.Achten Sie darauf, dass der Tabellentyp auf Native Tabelle eingestellt ist.
Geben Sie im Abschnitt Schema die Schemadefinition ein. Sie haben dazu folgende Möglichkeiten:
Sie können Schemainformationen manuell eingeben, indem Sie Als Text bearbeiten aktivieren und das Tabellenschema als JSON-Array eingeben. Geben Sie die folgenden Informationen ein:
[ { "name": "location", "type": "GEOGRAPHY", "mode": "REQUIRED" }, { "name": "average_temperature", "type": "INTEGER", "mode": "REQUIRED" }, { "name": "month", "type": "STRING", "mode": "REQUIRED" }, { "name": "inches_of_rain", "type": "NUMERIC" }, { "name": "is_current", "type": "BOOLEAN" }, { "name": "latest_measurement", "type": "DATE" } ]Geben Sie das Schema mithilfe von Feld hinzufügen manuell ein:

Behalten Sie für Partitions- und Clustereinstellungen den Standardwert
No partitioningbei.Übernehmen Sie im Abschnitt Erweiterte Optionen den Standardwert für Verschlüsselung, und zwar
Google-owned and managed key.Klicken Sie auf Tabelle erstellen.
bq
Verwenden Sie den Befehl bq mk, um ein leeres Dataset und eine Tabelle in diesem Dataset zu erstellen.
Führen Sie den folgenden Befehl aus, um ein Dataset mit dem durchschnittlichen globalen Wetter zu erstellen:
bq --location=LOCATION mk \
--dataset PROJECT_ID:average_weather
Ersetzen Sie Folgendes:
LOCATION: die Region, in der sich die Umgebung befindet.PROJECT_ID: Die Projekt-ID.
Führen Sie den folgenden Befehl aus, um eine leere Tabelle in diesem Dataset mit der Schemadefinition zu erstellen:
bq mk --table \
PROJECT_ID:average_weather.average_weather \
location:GEOGRAPHY,average_temperature:INTEGER,month:STRING,inches_of_rain:NUMERIC,is_current:BOOLEAN,latest_measurement:DATE
Nachdem die Tabelle erstellt wurde, können Sie die Ablaufzeit, die Beschreibung und die Labels der Tabelle aktualisieren. Ebenso können Sie die Schemadefinition ändern.
Python
Speichern Sie diesen Code als dataflowtemplateoperator_create_dataset_and_table_helper.py und aktualisieren Sie die Variablen darin, damit sie Ihr Projekt und Ihren Standort widerspiegeln. Führen Sie ihn dann mit dem folgenden Befehl aus:
python dataflowtemplateoperator_create_dataset_and_table_helper.py
Python
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Cloud Storage-Bucket erstellen
Erstellen Sie einen Bucket, in dem alle für den Workflow erforderlichen Dateien gespeichert werden. In der DAG, die Sie später in dieser Anleitung erstellen, wird auf die Dateien verwiesen, die Sie in diesen Speicher-Bucket hochladen. So erstellen Sie einen neuen Storage-Bucket:
Console
Öffnen Sie Cloud Storage in der Google Cloud Console.
Klicken Sie auf Bucket erstellen, um das Formular zum Erstellen eines Buckets zu öffnen.
Geben Sie die Bucket-Informationen ein und klicken Sie zum Ausführen der einzelnen Schritte jeweils auf Weiter:
Geben Sie einen global eindeutigen Namen für den Bucket an. In diesem Leitfaden wird
bucketNameals Beispiel verwendet.Wählen Sie Region als Standorttyp aus. Wählen Sie als Nächstes einen Standort für die Bucket-Daten aus.
Wählen Sie Standard als Standard-Speicherklasse für Ihre Daten aus.
Wählen Sie Einheitliche Zugriffssteuerung für den Zugriff auf Ihre Objekte aus.
Klicken Sie auf Fertig.
gcloud
Führen Sie den Befehl gcloud storage buckets create aus:
gcloud storage buckets create gs://bucketName/
Ersetzen Sie Folgendes:
bucketName: Der Name des Buckets, den Sie zuvor in diesem Leitfaden erstellt haben.
Codebeispiele
C#
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Go
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Java
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Python
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Ruby
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
BigQuery-Schema im JSON-Format für Ihre Ausgabetabelle erstellen
Erstellen Sie eine JSON-formatierte BigQuery-Schemadatei, die der zuvor erstellten Ausgabetabelle entspricht. Die Feldnamen, ‑typen und ‑modi müssen mit denen übereinstimmen, die zuvor in Ihrem BigQuery-Tabellenschema definiert wurden. In dieser Datei werden die Daten aus der Datei .txt in ein Format normalisiert, das mit Ihrem BigQuery-Schema kompatibel ist. Geben Sie dieser Datei den Namen jsonSchema.json.
{
"BigQuery Schema": [
{
"name": "location",
"type": "GEOGRAPHY",
"mode": "REQUIRED"
},
{
"name": "average_temperature",
"type": "INTEGER",
"mode": "REQUIRED"
},
{
"name": "month",
"type": "STRING",
"mode": "REQUIRED"
},
{
"name": "inches_of_rain",
"type": "NUMERIC"
},
{
"name": "is_current",
"type": "BOOLEAN"
},
{
"name": "latest_measurement",
"type": "DATE"
}]
}
JavaScript-Datei zum Formatieren Ihrer Daten erstellen
In dieser Datei definieren Sie Ihre benutzerdefinierte Funktion (User Defined Function, UDF), die die Logik für die Transformation der Textzeilen in Ihrer Eingabedatei bereitstellt. Beachten Sie, dass diese Funktion jede Textzeile in Ihrer Eingabedatei als eigenes Argument verwendet. Die Funktion wird also einmal für jede Zeile Ihrer Eingabedatei ausgeführt. Geben Sie dieser Datei den Namen transformCSVtoJSON.js.
Eingabedatei erstellen
Diese Datei enthält die Informationen, die Sie in Ihre BigQuery-Tabelle hochladen möchten. Kopieren Sie diese Datei lokal und nennen Sie sie inputFile.txt.
POINT(40.7128 74.006),45,'July',null,true,2020-02-16
POINT(41.8781 87.6298),23,'October',13,false,2015-02-13
POINT(48.8566 2.3522),80,'December',null,true,null
POINT(6.5244 3.3792),15,'March',14,true,null
Dateien in Ihren Bucket hochladen
Laden Sie die folgenden Dateien in den Cloud Storage-Bucket hoch, den Sie zuvor erstellt haben:
- JSON-formatiertes BigQuery-Schema (
.json) - Benutzerdefinierte JavaScript-Funktion (
transformCSVtoJSON.js) Die Eingabedatei für den zu verarbeitenden Text (
.txt)
Console
- Wechseln Sie in der Google Cloud Console unter „Cloud Storage“ zur Seite Buckets.
Klicken Sie in der Liste der Buckets auf Ihren Bucket.
Führen Sie im Tab „Objekte“ für den Bucket einen der folgenden Schritte aus:
Fügen Sie die gewünschten Dateien per Drag-and-drop von Ihrem Desktop oder Dateimanager in den Hauptbereich der Google Cloud Console ein.
Klicken Sie auf den Button Dateien hochladen, wählen Sie im angezeigten Dialogfeld die Dateien aus, die Sie hochladen möchten, und klicken Sie auf Öffnen.
gcloud
Führen Sie den Befehl gcloud storage cp aus:
gcloud storage cp OBJECT_LOCATION gs://bucketName
Ersetzen Sie Folgendes:
bucketName: der Name des Buckets, den Sie zuvor in dieser Anleitung erstellt haben.OBJECT_LOCATION: der lokale Pfad zu Ihrem Objekt. Beispiel:Desktop/transformCSVtoJSON.js
Codebeispiele
Python
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
Ruby
Richten Sie die Standardanmeldedaten für Anwendungen ein, um sich bei Managed Airflow zu authentifizieren. Weitere Informationen finden Sie unter Authentifizierung für eine lokale Entwicklungsumgebung einrichten.
DataflowTemplateOperator konfigurieren
Bevor Sie den DAG ausführen, legen Sie die folgenden Airflow-Variablen fest.
| Airflow-Variable | Wert |
|---|---|
project_id
|
Die Projekt-ID. Beispiel: example-project. |
gce_zone
|
Compute Engine-Zone, in der der Dataflow-Cluster erstellt werden muss. Beispiel: us-central1-a. Weitere Informationen zu gültigen Zonen finden Sie unter Regionen und Zonen. |
bucket_path
|
Der Speicherort des Cloud Storage-Bucket, den Sie zuvor erstellt haben. Beispiel: gs://example-bucket |
Nun verweisen Sie auf die Dateien, die Sie zuvor erstellt haben, um einen DAG zu erstellen, mit dem der Dataflow-Workflow gestartet wird. Kopieren Sie diesen DAG und speichern Sie ihn lokal als composer-dataflow-dag.py.
DAG in Cloud Storage hochladen
Laden Sie Ihren DAG in den Ordner /dags im Bucket Ihrer Umgebung hoch. Sobald der Upload erfolgreich abgeschlossen wurde, können Sie ihn in der Managed Airflow-Umgebung auf den Link DAGs-Ordner sehen.

Status der Aufgabe ansehen
- Rufen Sie die Airflow-Weboberfläche auf.
- Klicken Sie auf der Seite „DAGs“ auf den DAG-Namen, z. B.
composerDataflowDAG. - Klicken Sie auf der DAGs-Detailseite auf Grafikansicht.
Prüfen Sie den Status:
Failed: Die Aufgabe ist rot umrandet. Sie können auch den Mauszeiger über die Aufgabe halten und nach State: Failed suchen.Success: Die Aufgabe ist von einem grünen Kästchen umgeben. Sie können auch den Mauszeiger über die Aufgabe bewegen und nach Status: Erfolgreich suchen.
Nach einigen Minuten können Sie die Ergebnisse in Dataflow und BigQuery prüfen.
Job in Dataflow ansehen
Rufen Sie in der Google Cloud Console die Seite Dataflow auf.
Ihr Job heißt
dataflow_operator_transform_csv_to_bq. Am Ende des Namens ist mit einem Bindestrich eine eindeutige ID angehängt, z. B.:
Klicken Sie auf den Namen, um die Jobdetails aufzurufen.

Ergebnisse in BigQuery anzeigen
Rufen Sie in der Google Cloud Console die Seite BigQuery auf.
Sie können Abfragen mit Standard-SQL senden. Verwenden Sie die folgende Abfrage, um die Zeilen anzuzeigen, die Ihrer Tabelle hinzugefügt wurden:
SELECT * FROM projectId.average_weather.average_weather