【发布时间】:2021-09-24 05:06:32
【问题描述】:
我在 AWS 上的 EC2 机器上部署了气流调度程序和气流网络服务器。我使用这个气流调度程序来执行带有AwsBatchOperator 任务的 DAG。此任务执行 EC2 机器上的 python 脚本。这是 DAG 的代码:
from datetime import timedelta
from airflow import DAG
from airflow.utils.dates import days_ago
from airflow.providers.amazon.aws.operators.batch import AwsBatchOperator
default_args = {
'owner': 'admin',
'concurrency': 3,
'depends_on_past': True,
'email': ['airflow@example.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 1,
'retry_delay': timedelta(minutes=5),
'start_date': None,
'end_date': None,
'schedule_interval': None,
}
dag = DAG(
dag_id='my-dag',
default_args=default_args,
description='My DAG',
schedule_interval='00 03 * * *',
start_date=days_ago(1),
tags=['dev'],
)
task = AwsBatchOperator(
dag=dag,
job_name= 'my-job-name',
job_definition= 'arn:aws:batch:eu-central-1:XXXX:job-definition/my-job-name',
job_queue= 'arn:aws:batch:eu-central-1:XXXX:job-queue/my-job-name',
region_name= 'eu-central-1',
task_id= 'my-task-id',
overrides={
'command': ['python3', './my_python_script.py']
},
parameters= {}
)
python 脚本my_python_script.py 位于部署气流的EC2 机器上,在目录/home/ubuntu 中。
我在这个 python 脚本中出现了一个错误。我更正了它并将更正后的脚本推送到 EC2 机器上。但是,当我执行 DAG 时,我仍然收到由我更正的错字引起的错误。所以这是我的问题:
如何刷新我的 DAG 以确保它使用我的 EC2 机器上的脚本版本?
我试过的
- 点击 Airflow 网页界面上的“刷新”按钮刷新 DAG
- 等待气流调度程序自动刷新 DAG
- 删除 DAG 并等待刷新
- 在 EC2 机器上使用命令
python -m compileall重新编译 python 脚本
【问题讨论】:
-
我没有使用 AWS Batch,但也许您可以尝试使用相同的代码上传文件的副本,如
my_python_script_copy.py,并更改overrides参数中的名称以检查内容继续。
标签: python amazon-web-services airflow aws-batch