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 |
Pré-requisitos
Antes de começar, verifique se você tem:
Um cluster do AnalyticDB for MySQL nas edições Enterprise, Basic ou Data Lakehouse
Um grupo de recursos de job ou grupo de recursos interativos do Spark criado para o cluster
Python 3.7 ou superior
O endereço IP do servidor Airflow adicionado à lista de permissões do cluster
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
-
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 -
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:
ImportanteUse 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_typeMétodo de autenticação. Defina como
AKpara usar autenticação por par de AccessKey.access_key_idAccessKey ID do usuário RAM com acesso ao AnalyticDB for MySQL.
access_key_secretAccessKey secret do usuário RAM.
regionID 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>" } -
Crie um arquivo DAG chamado
spark_dags.py. O exemplo a seguir usa oAnalyticDBSparkSQLOperatorpara executar a consultaSHOW DATABASES:Parâmetros do DAG:
Parâmetro
Obrigatório
Descrição
dag_idSim
Nome do DAG.
default_argsSim
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_idSim
ID do job.
sqlSim
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 Copie o arquivo
spark_dags.pypara o diretóriodags_folderdefinido na configuração do Airflow.Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.
spark-submit
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
.
-
Instale o plug-in do Spark para o Airflow:
ImportanteA instalação do pacote
apache-airflow-providers-apache-sparkinstala automaticamente o PySpark. Para remover o PySpark, executepip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark -
Adicione o binário do spark-submit ao
PATHdo Airflow antes de iniciá-lo:ImportanteDefina o
PATHantes 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> -
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 Copie o arquivo
demo.pypara a pastadagsno diretório de instalação do Airflow.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.
-
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.
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.
No painel de navegação, escolha Cluster Management > Resource Management e clique na aba Resource Groups.
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.

-
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.
Na interface web do Airflow, acesse Admin > Connections.
-
Clique em
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
defaultpelo nome real do banco de dados e remova o sufixoresource_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
10000Extra
{"auth_mechanism": "CUSTOM"} -
Crie um arquivo DAG. O exemplo abaixo executa
show databasesdiariamente:Parâmetro
Obrigatório
Descrição
task_idSim
ID do job.
conn_idSim
Nome da conexão criada na etapa 4.
sqlSim
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 Na interface web do Airflow, clique em
ao lado do DAG para acioná-lo.
Agende jobs Spark JAR
Spark Airflow Operator
-
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 -
Crie uma conexão no Airflow. Na interface web do Airflow, acesse Admin > Connections e adicione uma conexão com o seguinte JSON:
ImportanteUse 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_typeMétodo de autenticação. Defina como
AK.access_key_idAccessKey ID do usuário RAM.
access_key_secretAccessKey secret do usuário RAM.
regionID 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>" } -
Crie um arquivo DAG chamado
spark_dags.py. O exemplo a seguir executa duas tarefas baseadas em JAR em sequência usando oAnalyticDBSparkBatchOperator:ImportanteArmazene 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_idSim
Nome do DAG.
default_argsSim
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_idSim
ID do job.
fileSim
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_nameObrigató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 Copie o arquivo
spark_dags.pypara o diretóriodags_folderdefinido na configuração do Airflow.Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.
spark-submit
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
.
-
Instale o plug-in do Spark para o Airflow:
ImportanteA instalação do pacote
apache-airflow-providers-apache-sparkinstala automaticamente o PySpark. Para remover o PySpark, executepip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark -
Adicione o binário do spark-submit ao
PATHdo Airflow antes de iniciá-lo:ImportanteDefina o
PATHantes 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> -
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_idSim
Nome do DAG.
default_argsSim
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_idSim
ID do job.
fileSim
Caminho absoluto para o arquivo principal da aplicação Spark. Deve estar armazenado no OSS.
class_nameObrigató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 Copie o arquivo
demo.pypara a pastadagsno diretório de instalação do Airflow.Acione o DAG pela interface web do Airflow. Para orientações, consulte o Tutorial do Airflow.