Todos os produtos
Search
Central de documentação

AnalyticDB:Use o Airflow para agendar jobs do Spark

Última atualização: Jun 27, 2026

Este guia explica como agendar jobs do Spark no AnalyticDB for MySQL com o Apache Airflow. O Airflow orquestra cargas de trabalho como grafos acíclicos direcionados (DAGs). Você pode conectar o Airflow ao AnalyticDB for MySQL de duas formas:

Método

Mais indicado para

Spark Airflow Operator

Integração mais estreita com o AnalyticDB for MySQL; usa autenticação por AccessKey; suporta jobs SQL e JAR

spark-submit

Provedor padrão do Apache Airflow para Spark; ideal se você já usa o pacote apache-airflow-providers-apache-spark

Pré-requisitos

Antes de começar, verifique se você tem:

Agende jobs do Spark SQL

O AnalyticDB for MySQL suporta Spark SQL nos modos batch e interativo. A configuração difere entre eles.

Modo batch

Spark Airflow Operator

  1. Instale o plug-in do Spark para o Airflow:

    pip install https://help-static-aliyun-doc.aliyuncs.com/file-manage-files/zh-CN/20230608/qvjf/adb_spark_airflow-0.0.1-py3-none-any.whl
  2. Crie uma conexão no Airflow. Na interface web do Airflow, acesse Admin > Connections e adicione uma conexão com o JSON abaixo como dados extras:

    Importante

    Use um usuário do Resource Access Management (RAM) com as permissões mínimas necessárias. Não use as credenciais da conta raiz da Alibaba Cloud.

    Parâmetro

    Descrição

    auth_type

    Método de autenticação. Defina como AK para usar autenticação por par de AccessKey.

    access_key_id

    AccessKey ID do usuário RAM com acesso ao AnalyticDB for MySQL.

    access_key_secret

    AccessKey secret do usuário RAM.

    region

    ID da região do cluster AnalyticDB for MySQL.

    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Crie um arquivo DAG chamado spark_dags.py. O exemplo a seguir usa o AnalyticDBSparkSQLOperator para executar a consulta SHOW DATABASES:

    Parâmetros do DAG:

    Parâmetro

    Obrigatório

    Descrição

    dag_id

    Sim

    Nome do DAG.

    default_args

    Sim

    Padrões no nível do cluster: cluster_id (ID do cluster), rg_name (nome do grupo de recursos de job), region (ID da região). Para mais informações, consulte Parâmetros do DAG.

    Parâmetros do AnalyticDBSparkSQLOperator:

    Parâmetro

    Obrigatório

    Descrição

    task_id

    Sim

    ID do job.

    sql

    Sim

    Instrução Spark SQL. Para mais informações, consulte Parâmetros do Airflow.

    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id="my_dag_name",
        default_args={"cluster_id": "<your_cluster_ID>", "rg_name": "<your_resource_group>", "region": "<your_region>"},
    ) as dag:
    
        spark_sql = AnalyticDBSparkSQLOperator(
            task_id="task2",
            sql="SHOW DATABASES;"
        )
    
        spark_sql
  4. Copie o arquivo spark_dags.py para o diretório dags_folder definido na configuração do Airflow.

  5. Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.

spark-submit

Nota

Você pode definir parâmetros específicos do AnalyticDB for MySQL (

clusterId

,

regionId

,

keyId

,

secretId

) no arquivo

conf/spark-defaults.conf

ou como parâmetros do Airflow. Para a lista completa, consulte

Parâmetros de configuração de aplicações Spark

.

  1. Instale o plug-in do Spark para o Airflow:

    Importante

    A instalação do pacote apache-airflow-providers-apache-spark instala automaticamente o PySpark. Para remover o PySpark, execute pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Baixe o pacote spark-submit e configure os parâmetros.

  3. Adicione o binário do spark-submit ao PATH do Airflow antes de iniciá-lo:

    Importante

    Defina o PATH antes de iniciar o Airflow. Se o Airflow iniciar sem o caminho do spark-submit, ele não localizará o comando durante a execução.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Crie um arquivo DAG chamado demo.py:

    from airflow.models import DAG
    from airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator
    from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
    from airflow.utils.dates import days_ago
    args = {
        'owner': 'Aliyun ADB Spark',
    }
    with DAG(
        dag_id='example_spark_operator',
        default_args=args,
        schedule_interval=None,
        start_date=days_ago(2),
        tags=['example'],
    ) as dag:
        adb_spark_conf = {
            "spark.driver.resourceSpec": "medium",
            "spark.executor.resourceSpec": "medium"
        }
        # Submit a Spark application from an OSS path
        submit_job = SparkSubmitOperator(
            conf=adb_spark_conf,
            application="oss://<bucket_name>/jar/pi.py",
            task_id="submit_job",
            verbose=True
        )
        # Run a Spark SQL query
        sql_job = SparkSqlOperator(
            conn_id="spark_default",
            sql="SELECT * FROM yourdb.yourtable",
            conf=",".join([k + "=" + v for k, v in adb_spark_conf.items()]),
            task_id="sql_job",
            verbose=True
        )
        submit_job >> sql_job
  5. Copie o arquivo demo.py para a pasta dags no diretório de instalação do Airflow.

  6. Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.

Modo interativo

No modo interativo, o Airflow conecta-se ao cluster AnalyticDB for MySQL por meio de um endpoint JDBC usando HiveServer2 Thrift.

  1. Obtenha o endpoint do grupo de recursos interativos do Spark. Clique em Apply for Endpoint ao lado de Public Endpoint para solicitar um endpoint público nos seguintes casos:

    • A ferramenta cliente usada para enviar jobs do Spark SQL está implantada em uma máquina local ou servidor externo.

    • A ferramenta cliente usada para enviar jobs do Spark SQL está implantada em uma instância ECS que não está na mesma VPC do cluster AnalyticDB for MySQL.

    1. Faça login no console do AnalyticDB for MySQL. No canto superior esquerdo, selecione uma região. No painel de navegação à esquerda, clique em Clusters e, em seguida, clique no ID do seu cluster.

    2. No painel de navegação, escolha Cluster Management > Resource Management e clique na aba Resource Groups.

    3. Localize o grupo de recursos desejado e clique em Details na coluna Actions. Copie o endpoint interno ou público. Você também pode copiar a string de conexão JDBC no campo Port. image

  2. Instale as dependências necessárias:

    pip install apache-airflow-providers-apache-hive "apache-airflow-providers-common-sql==1.21.0"

    Para detalhes sobre os pacotes, consulte apache-airflow-providers-apache-hive e apache-airflow-providers-common-sql.

  3. Na interface web do Airflow, acesse Admin > Connections.

  4. Clique em image para adicionar uma conexão. Configure os seguintes parâmetros:

    Parâmetro

    Descrição

    Connection Id

    Nome da conexão. Exemplo: adb_spark_cluster.

    Connection Type

    Selecione Hive Server 2 Thrift.

    Host

    Endpoint obtido na etapa 1. Substitua default pelo nome real do banco de dados e remova o sufixo resource_group=<resource group name>. Exemplo: jdbc:hive2://amv-t4naxpqk****sparkwho.ads.aliyuncs.com:10000/adb_demo.

    Schema

    Nome do banco de dados. Exemplo: adb_demo.

    Login

    Nome do grupo de recursos e conta do banco de dados no formato resource_group_name/database_account_name. Exemplo: spark_interactive_prod/spark_user.

    Password

    Senha da conta do banco de dados AnalyticDB for MySQL.

    Port

    10000

    Extra

    {"auth_mechanism": "CUSTOM"}

  5. Crie um arquivo DAG. O exemplo abaixo executa show databases diariamente:

    Parâmetro

    Obrigatório

    Descrição

    task_id

    Sim

    ID do job.

    conn_id

    Sim

    Nome da conexão criada na etapa 4.

    sql

    Sim

    Instrução Spark SQL.

    from airflow import DAG
    from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
    from datetime import datetime
    
    default_args = {
        'owner': 'airflow',
        'start_date': datetime(2025, 2, 10),
        'retries': 1,
    }
    
    dag = DAG(
        'adb_spark_sql_test',
        default_args=default_args,
        schedule_interval='@daily',
    )
    
    jdbc_query = SQLExecuteQueryOperator(
        task_id='execute_spark_sql_query',
        conn_id='adb_spark_cluster',  # Connection created in step 4
        sql='show databases',
        dag=dag
    )
    
    jdbc_query
  6. Na interface web do Airflow, clique em image ao lado do DAG para acioná-lo.

Agende jobs Spark JAR

Spark Airflow Operator

  1. Instale o plug-in do Spark para o Airflow:

    pip install https://help-static-aliyun-doc.aliyuncs.com/file-manage-files/zh-CN/20230608/qvjf/adb_spark_airflow-0.0.1-py3-none-any.whl
  2. Crie uma conexão no Airflow. Na interface web do Airflow, acesse Admin > Connections e adicione uma conexão com o seguinte JSON:

    Importante

    Use um usuário RAM com as permissões mínimas necessárias. Não use as credenciais da sua conta raiz. Para mais detalhes, consulte Contas e permissões.

    Parâmetro

    Descrição

    auth_type

    Método de autenticação. Defina como AK.

    access_key_id

    AccessKey ID do usuário RAM.

    access_key_secret

    AccessKey secret do usuário RAM.

    region

    ID da região do cluster AnalyticDB for MySQL.

    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Crie um arquivo DAG chamado spark_dags.py. O exemplo a seguir executa duas tarefas baseadas em JAR em sequência usando o AnalyticDBSparkBatchOperator:

    Importante

    Armazene todos os arquivos principais da aplicação Spark no Object Storage Service (OSS). O bucket do OSS e o cluster AnalyticDB for MySQL devem estar na mesma região.

    Parâmetros do DAG:

    Parâmetro

    Obrigatório

    Descrição

    dag_id

    Sim

    Nome do DAG.

    default_args

    Sim

    Padrões no nível do cluster: cluster_id, rg_name, region. Para mais informações, consulte Parâmetros do DAG.

    Parâmetros do AnalyticDBSparkBatchOperator:

    Parâmetro

    Obrigatório

    Descrição

    task_id

    Sim

    ID do job.

    file

    Sim

    Caminho absoluto para o arquivo principal da aplicação Spark — um pacote JAR (Java/Scala) ou arquivo de ponto de entrada Python. Deve estar armazenado no OSS.

    class_name

    Obrigatório para Java/Scala

    Classe de ponto de entrada. Aplicações Python não exigem este parâmetro. Para mais informações, consulte Parâmetros do AnalyticDBSparkBatchOperator.

    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id=DAG_ID,
        default_args={"cluster_id": "your cluster", "rg_name": "your resource group", "region": "your region"},
    ) as dag:
        spark_pi = AnalyticDBSparkBatchOperator(
            task_id="task1",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkPi",
        )
    
        spark_lr = AnalyticDBSparkBatchOperator(
            task_id="task2",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkLR",
        )
    
        spark_pi >> spark_lr
  4. Copie o arquivo spark_dags.py para o diretório dags_folder definido na configuração do Airflow.

  5. Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.

spark-submit

Nota

Você pode definir parâmetros específicos do AnalyticDB for MySQL (

clusterId

,

regionId

,

keyId

,

secretId

,

ossUploadPath

) no arquivo

conf/spark-defaults.conf

ou como parâmetros do Airflow. Para a lista completa, consulte

Parâmetros de configuração de aplicações Spark

.

  1. Instale o plug-in do Spark para o Airflow:

    Importante

    A instalação do pacote apache-airflow-providers-apache-spark instala automaticamente o PySpark. Para remover o PySpark, execute pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Baixe o pacote spark-submit e configure os parâmetros.

  3. Adicione o binário do spark-submit ao PATH do Airflow antes de iniciá-lo:

    Importante

    Defina o PATH antes de iniciar o Airflow. Se o Airflow iniciar sem o caminho do spark-submit, ele não localizará o comando durante a execução.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Crie um arquivo DAG chamado demo.py. O exemplo a seguir executa duas tarefas JAR em sequência:

    Parâmetros do DAG:

    Parâmetro

    Obrigatório

    Descrição

    dag_id

    Sim

    Nome do DAG.

    default_args

    Sim

    Padrões no nível do cluster: cluster_id, rg_name, region. Para mais informações, consulte Parâmetros do DAG.

    Parâmetros do AnalyticDBSparkBatchOperator:

    Parâmetro

    Obrigatório

    Descrição

    task_id

    Sim

    ID do job.

    file

    Sim

    Caminho absoluto para o arquivo principal da aplicação Spark. Deve estar armazenado no OSS.

    class_name

    Obrigatório para Java/Scala

    Classe de ponto de entrada. Aplicações Python não exigem este parâmetro. Para mais informações, consulte Parâmetros do AnalyticDBSparkBatchOperator.

    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id=DAG_ID,
        start_date=datetime(2021, 1, 1),
        schedule=None,
        default_args={"cluster_id": "your cluster", "rg_name": "your resource group", "region": "your region"},
        max_active_runs=1,
        catchup=False,
    ) as dag:
        spark_pi = AnalyticDBSparkBatchOperator(
            task_id="task1",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkPi",
        )
    
        spark_lr = AnalyticDBSparkBatchOperator(
            task_id="task2",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkLR",
        )
    
        spark_pi >> spark_lr
  5. Copie o arquivo demo.py para a pasta dags no diretório de instalação do Airflow.

  6. Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.