【发布时间】:2019-08-19 01:45:42
【问题描述】:
以下是气流 DAG 代码。当气流托管在本地和云作曲家时,它都能完美运行。但是,DAG 本身在 Composer UI 中是不可点击的。 我发现了一个类似的问题,并尝试了this question 中链接的公认答案。我的问题是类似的。
import airflow
from airflow import DAG
from airflow.operators.dummy_operator import DummyOperator
from airflow.operators.python_operator import PythonOperator
from airflow.operators.mysql_operator import MySqlOperator
from airflow.contrib.operators.dataproc_operator import DataprocClusterCreateOperator
from airflow.contrib.operators.dataproc_operator import DataprocClusterDeleteOperator
from airflow.contrib.operators.dataproc_operator import DataProcSparkOperator
from datetime import datetime, timedelta
import sys
#copy this package to dag directory in GCP composer bucket
from schemas.schemaValidator import loadSchema
from schemas.schemaValidator import sparkArgListToMap
#change these paths to point to GCP Composer data directory
## cluster config
clusterConfig= loadSchema("somePath/jobConfig/cluster.yaml","cluster")
##per job yaml config
autoLoanCsvToParquetConfig= loadSchema("somePath/jobConfig/job.yaml","job")
default_args= {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2019, 1, 1),
'retries': 1,
'retry_delay': timedelta(minutes=3)
}
dag= DAG('usr_job', default_args=default_args, schedule_interval=None)
t1= DummyOperator(task_id= "start", dag=dag)
t2= DataprocClusterCreateOperator(
task_id= "CreateCluster",
cluster_name= clusterConfig["cluster"]["cluster_name"],
project_id= clusterConfig["project_id"],
num_workers= clusterConfig["cluster"]["worker_config"]["num_instances"],
image_version= clusterConfig["cluster"]["dataproc_img"],
master_machine_type= clusterConfig["cluster"]["worker_config"]["machine_type"],
worker_machine_type= clusterConfig["cluster"]["worker_config"]["machine_type"],
zone= clusterConfig["region"],
dag=dag
)
t3= DataProcSparkOperator(
task_id= "csvToParquet",
main_class= autoLoanCsvToParquetConfig["job"]["main_class"],
arguments= autoLoanCsvToParquetConfig["job"]["args"],
cluster_name= clusterConfig["cluster"]["cluster_name"],
dataproc_spark_jars= autoLoanCsvToParquetConfig["job"]["jarPath"],
dataproc_spark_properties= sparkArgListToMap(autoLoanCsvToParquetConfig["spark_params"]),
dag=dag
)
t4= DataprocClusterDeleteOperator(
task_id= "deleteCluster",
cluster_name= clusterConfig["cluster"]["cluster_name"],
project_id= clusterConfig["project_id"],
dag= dag
)
t5= DummyOperator(task_id= "stop", dag=dag)
t1>>t2>>t3>>t4>>t5
UI 给出了这个错误 - "This DAG isn't available in the webserver DAG bag object. It shows up in this list because the scheduler marked it as active in the metadata database."
然而,当我在 Composer 上手动触发 DAG 时,我发现它通过日志文件成功运行。
【问题讨论】:
-
DAG 可能在您的集群中加载良好,但在 Web 服务器项目(与您的主项目不同)中加载不正常。你有来自 Stackdriver 的网络服务器日志可以发布吗?
-
我按照您的建议检查了日志,似乎没有任何异常。 Web 服务器和调度程序上没有任何类型的错误/警告。正如已经提到的,dag 可以顺利执行。只是如上面问题中的链接,dag 不可点击,我在网络服务器 UI 中看不到任何 dag 指标
-
您能否检查您的 Web 服务器日志并确保 GCS 同步从存储桶到 Web 服务器正常工作?其他 DAG 是否发生过这种情况?您的租户项目用户具有哪些 IAM 角色? (与您的 Web 服务器 URL 匹配的
*-tp@appspot.gserviceaccount.com用户) -
GCS 同步工作正常(否则,dag 不会执行。另外,我可以在 UI 中看到 dag,只是它不可点击。)我创建了 *-tp@appspot .gserviceaccount.com 服务帐户,并为简单起见为其分配了所有者角色。我尝试了另一个 DAG,似乎问题出在我共享的特定 DAG 上,但在任何地方都没有破损、错误/警告日志
标签: python airflow google-cloud-composer