Aktywowanie DAG-a w Node.JS i Google Cloud Functions

1. Wprowadzenie

Apache Airflow jest przeznaczony do regularnego uruchamiania DAG-ów, ale możesz też aktywować DAG-i w odpowiedzi na zdarzenia, takie jak zmiana w zasobniku Cloud Storage lub wiadomość przesłana do Cloud Pub/Sub. Aby to zrobić, możesz wywoływać DAG-i usługi zarządzanej Apache Airflow za pomocą Cloud Functions.

W przykładzie w tym module laboratorium prosty DAG jest uruchamiany za każdym razem, gdy w zasobniku Cloud Storage nastąpi zmiana. Ten DAG używa operatora BashOperator do uruchamiania polecenia bash, które wyświetla informacje o zmianach dotyczące tego, co zostało przesłane do zasobnika Cloud Storage.

Zanim rozpoczniesz to ćwiczenie, zapoznaj się z ćwiczeniami z programowania Wprowadzenie do zarządzanej usługi Apache Airflow i Wprowadzenie do Cloud Functions. Jeśli w ćwiczeniu z programowania Wprowadzenie do zarządzanej usługi Apache Airflow utworzysz środowisko zarządzanej usługi Apache Airflow, możesz go użyć w tym ćwiczeniu.

Co utworzysz

W tym ćwiczeniu:

  1. przesłać plik do Google Cloud Storage,
  2. Aktywowanie funkcji Google Cloud za pomocą środowiska wykonawczego Node.JS
  3. Ta funkcja uruchomi DAG w usłudze zarządzanej Apache Airflow w Google Cloud.
  4. Uruchamia proste polecenie bash, które wyświetla zmianę w zasobniku Google Cloud Storage.

1d3d3736624a923f.png

Czego się nauczysz

  • Jak aktywować graf DAG Apache Airflow za pomocą Google Cloud Functions i Node.js

Co będzie potrzebne

  • Konto GCP
  • Podstawowa znajomość języka JavaScript
  • Podstawowa wiedza o usłudze zarządzanej Apache Airflow/Airflow i Cloud Functions
  • wygodę korzystania z poleceń interfejsu wiersza poleceń,

2. Konfigurowanie GCP

Wybierz lub utwórz projekt

Wybierz lub utwórz projekt Google Cloud Platform. Jeśli tworzysz nowy projekt, wykonaj czynności opisane tutaj.

Zanotuj identyfikator projektu, który będzie Ci potrzebny w dalszych krokach.

Jeśli tworzysz nowy projekt, identyfikator projektu znajdziesz tuż pod nazwą projektu na stronie tworzenia.

Jeśli masz już utworzony projekt, jego identyfikator znajdziesz na stronie głównej konsoli na karcie Informacje o projekcie.

Włączanie interfejsów API

Włącz interfejsy Managed Apache Airflow API, Google Cloud Functions API, Cloud Identity API i Google Identity and Access Management (IAM) API.

Tworzenie środowiska zarządzanej usługi Apache Airflow

Utwórz środowisko zarządzanej usługi Apache Airflow o tej konfiguracji:

  • Nazwa: my-airflow-environment
  • Lokalizacja: dowolna lokalizacja geograficznie najbliższa użytkownikowi
  • Strefa: dowolna strefa w tym regionie

Wszystkie pozostałe konfiguracje mogą pozostać domyślne. U dołu kliknij „Utwórz”. Zanotuj nazwę i lokalizację środowiska zarządzanej usługi Apache Airflow – będą Ci potrzebne w kolejnych krokach.

Utworzenie zasobnika Cloud Storage

W projekcie utwórz zasobnik Cloud Storage o tej konfiguracji:

  • Nazwa: <your-project-id>
  • Domyślna klasa pamięci masowej: wiele regionów
  • Lokalizacja: dowolna lokalizacja, która jest geograficznie najbliżej regionu usługi zarządzanej Apache Airflow, z którego korzystasz.
  • Model kontroli dostępu: ustawienie uprawnień na poziomie obiektu i zasobnika

Gdy wszystko będzie gotowe, kliknij „Utwórz”. Zapamiętaj nazwę zasobnika Cloud Storage, przyda się w kolejnych krokach.

3. Konfigurowanie Google Cloud Functions (GCF)

Aby skonfigurować GCF, będziemy uruchamiać polecenia w Google Cloud Shell.

Z Google Cloud można korzystać zdalnie na laptopie za pomocą narzędzia wiersza poleceń gcloud. W tym ćwiczeniu użyjemy jednak Google Cloud Shell, czyli środowiska wiersza poleceń działającego w chmurze.

Ta maszyna wirtualna oparta na Debianie zawiera wszystkie potrzebne narzędzia dla programistów. Zawiera również stały katalog domowy o pojemności 5 GB i działa w Google Cloud, co znacznie zwiększa wydajność sieci i usprawnia proces uwierzytelniania. Oznacza to, że do ukończenia tego ćwiczenia potrzebujesz tylko przeglądarki (tak, działa ona na Chromebooku).

Aby aktywować Google Cloud Shell, w konsoli dewelopera kliknij przycisk w prawym górnym rogu (uzyskanie dostępu do środowiska i połączenie się z nim powinno zająć tylko kilka chwil):

Przyznawanie uprawnień do podpisywania obiektów blob kontu usługi Cloud Functions

Aby usługa GCF mogła uwierzytelniać się w Cloud IAP, czyli w proxy chroniącym serwer internetowy Airflow, musisz przyznać kontu usługi Appspot rolę Service Account Token Creator. Aby to zrobić, uruchom to polecenie w Cloud Shell, zastępując <your-project-id> nazwą projektu.

gcloud iam service-accounts add-iam-policy-binding \
<your-project-id>@appspot.gserviceaccount.com \
--member=serviceAccount:<your-project-id>@appspot.gserviceaccount.com \
--role=roles/iam.serviceAccountTokenCreator

Jeśli na przykład Twój projekt ma nazwę my-project, polecenie będzie wyglądać tak:

gcloud iam service-accounts add-iam-policy-binding \
my-project@appspot.gserviceaccount.com \
--member=serviceAccount:my-project@appspot.gserviceaccount.com \
--role=roles/iam.serviceAccountTokenCreator

Uzyskiwanie identyfikatora klienta

Aby utworzyć token do uwierzytelniania w Cloud IAP, funkcja wymaga identyfikatora klienta serwera proxy, który chroni serwer WWW Airflow. Interfejs Managed Apache Airflow API nie udostępnia tych informacji bezpośrednio. Zamiast tego wyślij nieuwierzytelnione żądanie do serwera internetowego Airflow i przechwyć identyfikator klienta z adresu URL przekierowania. Zrobimy to, uruchamiając plik Pythona za pomocą Cloud Shell, aby przechwycić identyfikator klienta.

Pobierz niezbędny kod z GitHuba, uruchamiając w Cloud Shell to polecenie:

cd
git clone https://github.com/GoogleCloudPlatform/python-docs-samples.git

Jeśli pojawi się błąd, ponieważ ten katalog już istnieje, zaktualizuj go do najnowszej wersji, uruchamiając to polecenie:

cd python-docs-samples/
git pull origin master

Przejdź do odpowiedniego katalogu, wpisując

cd python-docs-samples/composer/rest

Uruchom kod w Pythonie, aby uzyskać identyfikator klienta. Zastąp <your-project-id> nazwą projektu, <your-managed-airflow-location> lokalizacją utworzonego wcześniej środowiska Managed Apache Airflow, a <your-managed-airflow-environment> nazwą utworzonego wcześniej środowiska Managed Apache Airflow.

python3 get_client_id.py <your-project-id> <your-managed-airflow-location> <your-managed-airflow-environment>

Jeśli na przykład nazwa projektu to my-project, lokalizacja zarządzanej usługi Apache Airflow to us-central1, a nazwa środowiska to my-airflow-environment, polecenie będzie wyglądać tak:

python3 get_client_id.py my-project us-central1 my-airflow-environment

get_client_id.py wykonuje te czynności:

  • Uwierzytelnianie w Google Cloud
  • Wysyła nieuwierzytelnione żądanie HTTP do serwera WWW Airflow, aby uzyskać adres URI przekierowania.
  • wyodrębnia z tego przekierowania parametr zapytania client_id,
  • wydrukuje go, aby można było go użyć.

Identyfikator klienta zostanie wyświetlony w wierszu poleceń i będzie wyglądać mniej więcej tak:

12345678987654321-abc1def3ghi5jkl7mno8pqr0.apps.googleusercontent.com

4. Tworzenie funkcji

W Cloud Shell sklonuj repozytorium z niezbędnym przykładowym kodem, uruchamiając

cd
git clone https://github.com/GoogleCloudPlatform/nodejs-docs-samples.git

Przejdź do odpowiedniego katalogu i nie zamykaj Cloud Shell, aby wykonać kolejne kroki.

cd nodejs-docs-samples/composer/functions/composer-storage-trigger

Otwórz stronę Google Cloud Functions, klikając menu nawigacyjne, a następnie „Cloud Functions”.

U góry strony kliknij „UTWÓRZ FUNKCJĘ”.

Nazwij funkcję „my-function” i pozostaw domyślną ilość pamięci, czyli 256 MB.

Ustaw aktywator na „Cloud Storage”, pozostaw typ zdarzenia jako „Finalize/Create” (Zakończ/Utwórz) i przejdź do zasobnika utworzonego w kroku Tworzenie zasobnika Cloud Storage.

W sekcji Kod źródłowy pozostaw ustawienie „Edytor wbudowany”, a w sekcji Środowisko wykonawcze wybierz „Node.js 8”.

W Cloud Shell uruchom to polecenie. Spowoduje to otwarcie plików index.js i package.json w edytorze Cloud Shell.

cloudshell edit index.js package.json

Kliknij kartę package.json, skopiuj ten kod i wklej go w sekcji package.json w edytorze wbudowanym Cloud Functions.

Ustaw „Funkcję do wykonania” na triggerDag.

Kliknij kartę index.js, skopiuj kod i wklej go w sekcji index.js edytora wbudowanego Cloud Functions.

Zastąp PROJECT_ID identyfikatorem projektu, a CLIENT_ID identyfikatorem klienta zapisanym w kroku uzyskiwania identyfikatora klienta. Nie klikaj jeszcze „Utwórz” – musisz jeszcze wypełnić kilka pól.

W Cloud Shell uruchom to polecenie, zastępując <your-environment-name> nazwą środowiska Managed Apache Airflow, a <your-composer-region> regionem, w którym znajduje się środowisko Managed Apache Airflow.

gcloud composer environments describe <your-environment-name> --location <your-managed-airflow-region>

Jeśli na przykład Twoje środowisko ma nazwę my-airflow-environment i znajduje się w us-central1, polecenie będzie wyglądać tak:

gcloud composer environments describe my-airflow-environment --location us-central1

Dane wyjściowe powinny wyglądać mniej więcej tak:

config:
 airflowUri: https://abc123efghi456k-tp.appspot.com
 dagGcsPrefix: gs://narnia-north1-test-codelab-jklmno-bucket/dags
 gkeCluster: projects/a-project/zones/narnia-north1-b/clusters/narnia-north1-test-codelab-jklmno-gke
 nodeConfig:
   diskSizeGb: 100
   location: projects/a-project/zones/narnia-north1-b
   machineType: projects/a-project/zones/narnia-north1-b/machineTypes/n1-standard-1
   network: projects/a-project/global/networks/default
   oauthScopes:
   - https://www.googleapis.com/auth/cloud-platform
   serviceAccount: 987665432-compute@developer.gserviceaccount.com
 nodeCount: 3
 softwareConfig:
   imageVersion: composer-1.7.0-airflow-1.10.0
   pythonVersion: '2'
createTime: '2019-05-29T09:41:27.919Z'
name: projects/a-project/locations/narnia-north1/environments/my-airflow-environment
state: RUNNING
updateTime: '2019-05-29T09:56:29.969Z'
uuid: 123456-7890-9876-543-210123456

W tych danych wyjściowych poszukaj zmiennej o nazwie airflowUri. W kodzie index.js zmień WEBSERVER_ID na identyfikator serwera internetowego Airflow – jest to część zmiennej airflowUri, która na końcu ma „-tp”, np. abc123efghi456k-tp.

Kliknij link „Więcej” w menu, a potem wybierz region, który znajduje się najbliżej Ciebie.

Zaznacz „Ponów próbę w przypadku niepowodzenia”

Aby utworzyć funkcję w Cloud Functions, kliknij „Utwórz”.

Krokowe wykonywanie kodu

Kod skopiowany z pliku index.js będzie wyglądać mniej więcej tak:

// [START composer_trigger]
'use strict';

const fetch = require('node-fetch');
const FormData = require('form-data');

/**
 * Triggered from a message on a Cloud Storage bucket.
 *
 * IAP authorization based on:
 * https://stackoverflow.com/questions/45787676/how-to-authenticate-google-cloud-functions-for-access-to-secure-app-engine-endpo
 * and
 * https://cloud.google.com/iap/docs/authentication-howto
 *
 * @param {!Object} data The Cloud Functions event data.
 * @returns {Promise}
 */
exports.triggerDag = async data => {
  // Fill in your Managed Apache Airflow environment information here.

  // The project that holds your function
  const PROJECT_ID = 'your-project-id';
  // Navigate to your webserver's login page and get this from the URL
  const CLIENT_ID = 'your-iap-client-id';
  // This should be part of your webserver's URL:
  // {tenant-project-id}.appspot.com
  const WEBSERVER_ID = 'your-tenant-project-id';
  // The name of the DAG you wish to trigger
  const DAG_NAME = 'composer_sample_trigger_response_dag';

  // Other constants
  const WEBSERVER_URL = `https://${WEBSERVER_ID}.appspot.com/api/experimental/dags/${DAG_NAME}/dag_runs`;
  const USER_AGENT = 'gcf-event-trigger';
  const BODY = {conf: JSON.stringify(data)};

  // Make the request
  try {
    const iap = await authorizeIap(CLIENT_ID, PROJECT_ID, USER_AGENT);

    return makeIapPostRequest(
      WEBSERVER_URL,
      BODY,
      iap.idToken,
      USER_AGENT,
      iap.jwt
    );
  } catch (err) {
    throw new Error(err);
  }
};

/**
 * @param {string} clientId The client id associated with the Managed Apache Airflow webserver application.
 * @param {string} projectId The id for the project containing the Cloud Function.
 * @param {string} userAgent The user agent string which will be provided with the webserver request.
 */
const authorizeIap = async (clientId, projectId, userAgent) => {
  const SERVICE_ACCOUNT = `${projectId}@appspot.gserviceaccount.com`;
  const JWT_HEADER = Buffer.from(
    JSON.stringify({alg: 'RS256', typ: 'JWT'})
  ).toString('base64');

  let jwt = '';
  let jwtClaimset = '';

  // Obtain an Oauth2 access token for the appspot service account
  const res = await fetch(
    `http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/${SERVICE_ACCOUNT}/token`,
    {
      headers: {'User-Agent': userAgent, 'Metadata-Flavor': 'Google'},
    }
  );
  const tokenResponse = await res.json();
  if (tokenResponse.error) {
    return Promise.reject(tokenResponse.error);
  }

  const accessToken = tokenResponse.access_token;
  const iat = Math.floor(new Date().getTime() / 1000);
  const claims = {
    iss: SERVICE_ACCOUNT,
    aud: 'https://www.googleapis.com/oauth2/v4/token',
    iat: iat,
    exp: iat + 60,
    target_audience: clientId,
  };
  jwtClaimset = Buffer.from(JSON.stringify(claims)).toString('base64');
  const toSign = [JWT_HEADER, jwtClaimset].join('.');

  const blob = await fetch(
    `https://iam.googleapis.com/v1/projects/${projectId}/serviceAccounts/${SERVICE_ACCOUNT}:signBlob`,
    {
      method: 'POST',
      body: JSON.stringify({
        bytesToSign: Buffer.from(toSign).toString('base64'),
      }),
      headers: {
        'User-Agent': userAgent,
        Authorization: `Bearer ${accessToken}`,
      },
    }
  );
  const blobJson = await blob.json();
  if (blobJson.error) {
    return Promise.reject(blobJson.error);
  }

  // Request service account signature on header and claimset
  const jwtSignature = blobJson.signature;
  jwt = [JWT_HEADER, jwtClaimset, jwtSignature].join('.');
  const form = new FormData();
  form.append('grant_type', 'urn:ietf:params:oauth:grant-type:jwt-bearer');
  form.append('assertion', jwt);

  const token = await fetch('https://www.googleapis.com/oauth2/v4/token', {
    method: 'POST',
    body: form,
  });
  const tokenJson = await token.json();
  if (tokenJson.error) {
    return Promise.reject(tokenJson.error);
  }

  return {
    jwt: jwt,
    idToken: tokenJson.id_token,
  };
};

/**
 * @param {string} url The url that the post request targets.
 * @param {string} body The body of the post request.
 * @param {string} idToken Bearer token used to authorize the iap request.
 * @param {string} userAgent The user agent to identify the requester.
 */
const makeIapPostRequest = async (url, body, idToken, userAgent) => {
  const res = await fetch(url, {
    method: 'POST',
    headers: {
      'User-Agent': userAgent,
      Authorization: `Bearer ${idToken}`,
    },
    body: JSON.stringify(body),
  });

  if (!res.ok) {
    const err = await res.text();
    throw new Error(err);
  }
};
// [END composer_trigger]

Przyjrzyjmy się, co się dzieje. Mamy tu 3 funkcje: triggerDag, authorizeIap i makeIapPostRequest.

triggerDag to funkcja, która jest wywoływana, gdy przesyłamy coś do wyznaczonego zasobnika Cloud Storage. W tym miejscu konfigurujemy ważne zmienne używane w innych żądaniach, takie jak PROJECT_ID, CLIENT_ID, WEBSERVER_ID i DAG_NAME. Wywołuje funkcje authorizeIap i makeIapPostRequest.

exports.triggerDag = async data => {
  // Fill in your Managed Apache Airflow environment information here.

  // The project that holds your function
  const PROJECT_ID = 'your-project-id';
  // Navigate to your webserver's login page and get this from the URL
  const CLIENT_ID = 'your-iap-client-id';
  // This should be part of your webserver's URL:
  // {tenant-project-id}.appspot.com
  const WEBSERVER_ID = 'your-tenant-project-id';
  // The name of the DAG you wish to trigger
  const DAG_NAME = 'composer_sample_trigger_response_dag';

  // Other constants
  const WEBSERVER_URL = `https://${WEBSERVER_ID}.appspot.com/api/experimental/dags/${DAG_NAME}/dag_runs`;
  const USER_AGENT = 'gcf-event-trigger';
  const BODY = {conf: JSON.stringify(data)};

  // Make the request
  try {
    const iap = await authorizeIap(CLIENT_ID, PROJECT_ID, USER_AGENT);

    return makeIapPostRequest(
      WEBSERVER_URL,
      BODY,
      iap.idToken,
      USER_AGENT,
      iap.jwt
    );
  } catch (err) {
    throw new Error(err);
  }
};

authorizeIap wysyła żądanie do serwera proxy, który chroni serwer WWW Airflow, używając konta usługi i „wymieniając” token JWT na token tożsamości, który będzie używany do uwierzytelniania authorizeIap.makeIapPostRequest

const authorizeIap = async (clientId, projectId, userAgent) => {
  const SERVICE_ACCOUNT = `${projectId}@appspot.gserviceaccount.com`;
  const JWT_HEADER = Buffer.from(
    JSON.stringify({alg: 'RS256', typ: 'JWT'})
  ).toString('base64');

  let jwt = '';
  let jwtClaimset = '';

  // Obtain an Oauth2 access token for the appspot service account
  const res = await fetch(
    `http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/${SERVICE_ACCOUNT}/token`,
    {
      headers: {'User-Agent': userAgent, 'Metadata-Flavor': 'Google'},
    }
  );
  const tokenResponse = await res.json();
  if (tokenResponse.error) {
    return Promise.reject(tokenResponse.error);
  }

  const accessToken = tokenResponse.access_token;
  const iat = Math.floor(new Date().getTime() / 1000);
  const claims = {
    iss: SERVICE_ACCOUNT,
    aud: 'https://www.googleapis.com/oauth2/v4/token',
    iat: iat,
    exp: iat + 60,
    target_audience: clientId,
  };
  jwtClaimset = Buffer.from(JSON.stringify(claims)).toString('base64');
  const toSign = [JWT_HEADER, jwtClaimset].join('.');

  const blob = await fetch(
    `https://iam.googleapis.com/v1/projects/${projectId}/serviceAccounts/${SERVICE_ACCOUNT}:signBlob`,
    {
      method: 'POST',
      body: JSON.stringify({
        bytesToSign: Buffer.from(toSign).toString('base64'),
      }),
      headers: {
        'User-Agent': userAgent,
        Authorization: `Bearer ${accessToken}`,
      },
    }
  );
  const blobJson = await blob.json();
  if (blobJson.error) {
    return Promise.reject(blobJson.error);
  }

  // Request service account signature on header and claimset
  const jwtSignature = blobJson.signature;
  jwt = [JWT_HEADER, jwtClaimset, jwtSignature].join('.');
  const form = new FormData();
  form.append('grant_type', 'urn:ietf:params:oauth:grant-type:jwt-bearer');
  form.append('assertion', jwt);

  const token = await fetch('https://www.googleapis.com/oauth2/v4/token', {
    method: 'POST',
    body: form,
  });
  const tokenJson = await token.json();
  if (tokenJson.error) {
    return Promise.reject(tokenJson.error);
  }

  return {
    jwt: jwt,
    idToken: tokenJson.id_token,
  };
};

makeIapPostRequest wysyła wywołanie do serwera WWW Airflow, aby wywołać composer_sample_trigger_response_dag. Nazwa DAG jest osadzona w adresie URL serwera WWW Airflow przekazywanym z parametrem url, a idToken to token uzyskany w żądaniu authorizeIap.

const makeIapPostRequest = async (url, body, idToken, userAgent) => {
  const res = await fetch(url, {
    method: 'POST',
    headers: {
      'User-Agent': userAgent,
      Authorization: `Bearer ${idToken}`,
    },
    body: JSON.stringify(body),
  });

  if (!res.ok) {
    const err = await res.text();
    throw new Error(err);
  }
};

5. Konfigurowanie DAG-a

W Cloud Shell przejdź do katalogu z przykładowymi przepływami pracy. Jest on częścią pakietu python-docs-samples pobranego z GitHuba w kroku Pobieranie identyfikatora klienta.

cd
cd python-docs-samples/composer/workflows

Prześlij DAG-a do usługi Managed Apache Airflow

Prześlij przykładowy DAG do zasobnika pamięci DAG środowiska Managed Apache Airflow za pomocą tego polecenia, gdzie <environment_name> to nazwa środowiska Managed Apache Airflow, a <location> to nazwa regionu, w którym się ono znajduje. trigger_response_dag.py to DAG, z którym będziemy pracować.

gcloud composer environments storage dags import \
    --environment <environment_name> \
    --location <location> \
    --source trigger_response_dag.py

Jeśli na przykład Twoje środowisko Managed Apache Airflow ma nazwę my-airflow-environment i znajduje się w us-central1, polecenie będzie wyglądać tak:

gcloud composer environments storage dags import \
    --environment my-airflow-environment \
    --location us-central1 \
    --source trigger_response_dag.py

Przechodzenie przez DAG-a

Kod DAG w trigger_response.py wygląda tak:

import datetime
import airflow
from airflow.operators import bash_operator


default_args = {
    'owner': 'Managed Apache Airflow Example',
    'depends_on_past': False,
    'email': [''],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': datetime.timedelta(minutes=5),
    'start_date': datetime.datetime(2017, 1, 1),
}

with airflow.DAG(
        'composer_sample_trigger_response_dag',
        default_args=default_args,
        # Not scheduled, trigger only
        schedule_interval=None) as dag:

    # Print the dag_run's configuration, which includes information about the
    # Cloud Storage object change.
    print_gcs_info = bash_operator.BashOperator(
        task_id='print_gcs_info', bash_command='echo {{ dag_run.conf }}')

Sekcja default_args zawiera argumenty domyślne wymagane przez model BaseOperator w Apache Airflow. Tę sekcję z tymi parametrami zobaczysz w każdym DAG-u Apache Airflow. Wartość owner jest obecnie ustawiona na Managed Apache Airflow Example, ale możesz ją zmienić na swoje imię i nazwisko. depends_on_past pokazuje, że ten DAG nie jest zależny od żadnych poprzednich DAG-ów. Trzy sekcje e-maili, email, email_on_failure i email_on_retry, są ustawione tak, aby nie przychodziły żadne powiadomienia e-mail na podstawie stanu tego DAG-u. DAG będzie ponawiany tylko raz, ponieważ wartość retries jest ustawiona na 1, i będzie to robić po 5 minutach, zgodnie z ustawieniem retry_delay. Wartość start_date zwykle określa, kiedy DAG powinien być uruchamiany, w połączeniu z jego schedule_interval (ustawionym później), ale w przypadku tego DAG-u nie ma znaczenia. Jest ustawiona na 1 stycznia 2017 roku, ale może być ustawiona na dowolną datę w przeszłości.

default_args = {
    'owner': 'Managed Apache Airflow Example',
    'depends_on_past': False,
    'email': [''],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': datetime.timedelta(minutes=5),
    'start_date': datetime.datetime(2017, 1, 1),
}

Sekcja with airflow.DAG konfiguruje DAG, który zostanie uruchomiony. Zostanie ona uruchomiona z identyfikatorem zadania composer_sample_trigger_response_dag, argumentami domyślnymi z sekcji default_args i co najważniejsze, z wartością schedule_interval równą None. Wartość schedule_interval jest ustawiona na None, ponieważ aktywujemy ten konkretny graf DAG za pomocą funkcji Cloud Function. Dlatego start_date w default_args nie ma znaczenia.

Po uruchomieniu DAG wyświetli swoją konfigurację zgodnie z instrukcjami w zmiennej print_gcs_info.

with airflow.DAG(
        'composer_sample_trigger_response_dag',
        default_args=default_args,
        # Not scheduled, trigger only
        schedule_interval=None) as dag:

    # Print the dag_run's configuration, which includes information about the
    # Cloud Storage object change.
    print_gcs_info = bash_operator.BashOperator(
        task_id='print_gcs_info', bash_command='echo {{ dag_run.conf }}')

6. Testowanie funkcji

Otwórz środowisko Managed Apache Airflow i w wierszu z nazwą środowiska kliknij link Airflow.

Otwórz composer_sample_trigger_response_dag, klikając jego nazwę. Obecnie nie ma żadnych dowodów na uruchomienie DAG-u, ponieważ jeszcze go nie uruchomiliśmy.Jeśli ten DAG nie jest widoczny lub nie można go kliknąć, poczekaj minutę i odśwież stronę.

Otwórz osobną kartę i prześlij dowolny plik do utworzonego wcześniej zasobnika Cloud Storage, który został określony jako wyzwalacz funkcji Cloud. Możesz to zrobić za pomocą konsoli lub polecenia gsutil.

Wróć do karty z interfejsem Airflow i kliknij Widok grafu.

Kliknij zadanie print_gcs_info, które powinno być otoczone zieloną linią.

W prawym górnym rogu menu kliknij „Wyświetl dziennik”.

W logach zobaczysz informacje o pliku przesłanym do zasobnika Cloud Storage.

Gratulacje! Właśnie udało Ci się aktywować graf DAG Airflow za pomocą Node.js i Google Cloud Functions.

7. Czyszczenie

Aby uniknąć obciążenia konta Google Cloud Platform opłatami za zasoby użyte w tym przewodniku:

  1. (Opcjonalnie) Aby zapisać dane, pobierz je z zasobnika Cloud Storage w środowisku Managed Apache Airflow i z utworzonego przez siebie zasobnika pamięci na potrzeby tego przewodnika.
  2. Usuń zasobnik Cloud Storage dla utworzonego środowiska.
  3. Usuń środowisko usługi zarządzanej Apache Airflow. Pamiętaj, że usunięcie środowiska nie powoduje usunięcia zasobnika pamięci środowiska.
  4. (Opcjonalnie) W przypadku bezserwerowego przetwarzania danych pierwsze 2 miliony wywołań miesięcznie są bezpłatne, a gdy skalujesz funkcję do zera, nie ponosisz żadnych opłat (więcej informacji znajdziesz na stronie cennika). Jeśli jednak chcesz usunąć funkcję w Cloud Functions, kliknij „USUŃ” w prawym górnym rogu strony z informacjami o funkcji.

4fe11e1b41b32ba2.png

Opcjonalnie możesz też usunąć projekt:

  1. W konsoli GCP otwórz stronę Projekty.
  2. Z listy projektów wybierz projekt, który chcesz usunąć, i kliknij Usuń.
  3. W polu wpisz identyfikator projektu i kliknij Wyłącz, aby usunąć projekt.