Airflow DAG部署后无执行:GCS无格式文件转CSV故障求助
问题排查与解决方案
一、基础执行触发问题
你的DAG设置了schedule_interval=None,这意味着它不会自动调度执行,必须手动触发:
- 登录Airflow UI,找到
dag_cartolas_gcs_to_csvDAG - 点击右侧的"触发DAG"按钮(播放图标)手动启动一次运行
二、DAG加载与权限验证
确认DAG是否被Airflow识别
- 检查Airflow UI的"DAGs"页面,确认
dag_cartolas_gcs_to_csv是否在列表中,状态是否为"已启用"(开关为绿色) - 若未出现,排查:
- DAG文件是否放在Airflow指定的
dags目录下 - 文件权限是否允许Airflow进程读取
- 查看Airflow scheduler日志,排查是否存在语法错误或加载失败信息
- DAG文件是否放在Airflow指定的
- 检查Airflow UI的"DAGs"页面,确认
GCS权限验证
- 确保Airflow使用的服务账号拥有
cartolas-ahorros-test存储桶的storage.objects.list、storage.objects.get和storage.objects.create权限 - 可在Airflow服务器上执行以下命令测试权限(替换服务账号密钥路径):
gcloud auth activate-service-account --key-file=/path/to/your-key.json gsutil ls gs://cartolas-ahorros-test/CARTOLAS.*
- 确保Airflow使用的服务账号拥有
三、代码逻辑优化
1. 过滤已处理的CSV文件
当前代码会列出所有以CARTOLAS.开头的对象,包括已生成的.csv文件,需添加过滤规则:
list_objects = GCSListObjectsOperator( task_id='list_objects', bucket='cartolas-ahorros-test', prefix='CARTOLAS.', exclude_pattern='*.csv', # 排除已处理的CSV文件 dag=dag, )
2. 支持多文件处理
当前代码仅处理列表中的第一个对象,若需批量处理所有符合条件的文件,修改download_task:
def download_task(**kwargs): bucket_name = 'cartolas-ahorros-test' ti = kwargs['ti'] objects = ti.xcom_pull(task_ids='list_objects') if not objects: raise ValueError("未找到前缀为'CARTOLAS.'的对象") local_paths = [] for object_name in objects: local_path = '/tmp/' + object_name download_file_from_gcs(bucket_name, object_name, local_path) local_paths.append(local_path) return local_paths
同时需修改后续process_task和upload_task,适配列表类型的输入参数。
3. 添加日志输出
在Python函数中添加日志,便于排查执行细节:
import logging logger = logging.getLogger(__name__) def procesar_archivo(input_path, output_path): logger.info(f"开始处理文件: {input_path}") # 原文件处理逻辑... logger.info(f"文件处理完成,输出路径: {output_path}") def download_task(**kwargs): logger.info("开始下载GCS文件") # 原下载逻辑... logger.info(f"下载完成,本地路径列表: {local_paths}")
4. 清理临时文件
处理完成后删除本地临时文件,避免占用磁盘空间:
def upload_task(**kwargs): ti = kwargs['ti'] output_paths = ti.xcom_pull(task_ids='process_file') bucket_name = 'cartolas-ahorros-test' for output_path in output_paths: object_name = os.path.basename(output_path) upload_file_to_gcs(bucket_name, object_name, output_path) # 删除临时文件 os.remove(output_path) os.remove(output_path.replace('.csv', '')) logger.info(f"已清理临时文件: {output_path}")
四、日志排查技巧
- 手动触发后若无日志,检查Airflow任务日志路径(默认
$AIRFLOW_HOME/logs/dag_id/task_id/execution_date) - 查看scheduler日志,确认DAG是否被正确调度
- 查看worker日志,确认任务是否被分配执行
内容的提问来源于stack exchange,提问作者Nelson Romero
相关产品推荐
相关产品推荐

