MaxCompute vous permet de planifier des tâches via des interfaces Python en utilisant Apache Airflow. Cette rubrique explique comment utiliser les opérateurs Python d'Apache Airflow pour orchestrer des tâches MaxCompute.
Contexte
Développé initialement par Airbnb, Apache Airflow est un outil open source écrit en Python et dédié à la planification de tâches. Il s'appuie sur un graphe orienté acyclique (DAG) pour définir un ensemble de tâches interdépendantes et les ordonnancer selon leurs relations de dépendance. Apache Airflow permet également de définir des sous-tâches via des interfaces Python et propose divers opérateurs adaptés à vos besoins métier. Pour plus d'informations, consultez la documentation Apache Airflow.
Prérequis
Avant d'utiliser Apache Airflow pour planifier des tâches MaxCompute, assurez-vous que les conditions suivantes sont remplies :
-
Apache Airflow est installé et démarré.
Pour plus de détails, consultez le guide Quick Start.
Cette procédure utilise la version 1.10.7 d'Apache Airflow.
Étape 1 : Rédiger un script Python de planification et l'enregistrer dans le répertoire racine d'Apache Airflow
Rédigez un script Python contenant la logique de planification complète ainsi que le nom de la tâche à orchestrer, puis enregistrez-le au format .py. Dans cet exemple, nous créons un fichier nommé Airflow_MC.py contenant le code suivant :
# -*- coding: UTF-8 -*-
import sys
import os
from odps import ODPS
from odps import options
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime, timedelta
from configparser import ConfigParser
import time
reload(sys)
sys.setdefaultencoding('utf8')
# Change the default encoding format.
# MaxCompute parameter settings
options.sql.settings = {'options.tunnel.limit_instance_tunnel': False, 'odps.sql.allow.fullscan': True}
cfg = ConfigParser()
cfg.read("odps.ini")
print(cfg.items())
# Replace the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable with the AccessKey ID of the user account.
# Replace the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable with the AccessKey secret of the user account.
# We recommend that you do not directly use the strings of your AccessKey ID and AccessKey secret.
odps = ODPS(cfg.get("odps",os.getenv('ALIBABA_CLOUD_ACCESS_KEY_ID')),cfg.get("odps",os.getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')),cfg.get("odps","project"),cfg.get("odps","endpoint"))
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'retry_delay': timedelta(minutes=5),
'start_date':datetime(2020,1,15)
# 'email': ['airflow@example.com'],
# 'email_on_failure': False,
# 'email_on_retry': False,
# 'retries': 1,
# 'queue': 'bash_queue',
# 'pool': 'backfill',
# 'priority_weight': 10,
# 'end_date': datetime(2016, 1, 1),
}
# Scheduling workflow
dag = DAG(
'Airflow_MC', default_args=default_args, schedule_interval=timedelta(seconds=30))
def read_sql(sqlfile):
with io.open(sqlfile, encoding='utf-8', mode='r') as f:
sql=f.read()
f.closed
return sql
# Job scheduling
def get_time():
print 'Current time {}'.format(time.time())
return time.time()
# Job scheduling
def mc_job ():
project = odps.get_project() # Obtain information of the default project.
instance=odps.run_sql("select * from long_chinese;")
print(instance.get_logview_address())
instance.wait_for_success()
with instance.open_reader() as reader:
count = reader.count
print("Number of data records in the table: {}".format(count))
for record in reader:
print record
return count
t1 = PythonOperator (
task_id = 'get_time' ,
provide_context = False ,
python_callable = get_time,
dag = dag )
t2 = PythonOperator (
task_id = 'mc_job' ,
provide_context = False ,
python_callable = mc_job ,
dag = dag )
t2.set_upstream(t1)
Étape 2 : Soumettre le script de planification
-
Dans la fenêtre de ligne de commande, exécutez la commande suivante pour soumettre le script Python rédigé lors de l'Étape 1.
python Airflow_MC.py -
Dans la fenêtre de ligne de commande, exécutez les commandes ci-dessous pour générer le workflow de planification et lancer une tâche de test.
# print the list of active DAGs airflow list_dags # prints the list of tasks the "tutorial" dag_id airflow list_tasks Airflow_MC # prints the hierarchy of tasks in the tutorial DAG airflow list_tasks Airflow_MC --tree # Run a test job. airflow test Airflow_MC get_time 2010-01-16 airflow test Airflow_MC mc_job 2010-01-16
Étape 3 : Exécuter une tâche
Connectez-vous à l'interface web d'Apache Airflow. Sur la page DAGs, repérez le workflow soumis, puis cliquez sur l'icône
située dans la colonne Links pour lancer l'exécution.

Étape 4 : Consulter le résultat d'exécution
Cliquez sur le nom de la tâche pour afficher son workflow dans l'onglet Graph View. Sélectionnez ensuite une tâche spécifique au sein du flux, telle que mc_job. Dans la boîte de dialogue qui s'affiche, cliquez sur View Log pour visualiser le journal d'exécution.
