You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Airflow中捕获Scrapy爬虫异常并标记任务失败

问题

通过Airflow的PythonOperator调用Scrapy爬虫时,即使爬虫运行中出现GCS凭证异常(google.auth.exceptions.DefaultCredentialsError),Airflow任务仍被标记为SUCCESS,需要捕获该异常并将任务标记为FAILED。

当前代码

DAG定义

with DAG(
    'MySpider',
    default_args=default_args,
    schedule_interval=None) as dag:

    t1 = python_task = PythonOperator(
        task_id="crawler_task",
        python_callable=run_crawler,
        op_kwargs=dag_kwargs
    )

run_crawler方法

def run_crawler(**kwargs):
    project_settings = set_project_settings({
        'FEEDS': {
            f'{kwargs["bucket"]}%(time)s.{kwargs["format"]}': {
                'format': kwargs["format"],
                'encoding': 'utf8',
                'store_empty': kwargs["store_empty"]
            }
        }
    })
    
    print("Project settings: ")
    pprint(project_settings.attributes.items())

    set_connection("airflow", kwargs["gcs_connection_id"])
    
    process = CrawlerProcess(project_settings)
    process.crawl(spider.MySpider)

    print("Starting crawler...")
    process.start()

异常日志片段

google.auth.exceptions.DefaultCredentialsError: The file /tmp/file_my_credentials.json does not have a valid type. Type is None, expected one of ('authorized_user', 'service_account', 'external_account', 'external_account_authorized_user', 'impersonated_service_account', 'gdch_service_account').

{logging_mixin.py:115} WARNING - [scrapy.statscollectors] INFO: Dumping Scrapy stats:
{
    ...
    'feedexport/failed_count/GCSFeedStorage': 1,
    'log_count/ERROR': 1,
    ...
}
[2032-13-13, 09:04:28 UTC] {taskinstance.py:1408} INFO - Marking task as SUCCESS. dag_id=MySpider, task_id=crawler_task, execution_date=2032-13-13, start_date=2032-13-13, end_date=2032-13-13
解决方案

原因分析

Scrapy的CrawlerProcess.start()会启动Twisted reactor,运行过程中的异常会被Scrapy内部日志系统捕获,但不会向上冒泡到Airflow的PythonOperator上下文,导致Airflow无法感知异常,默认标记任务为SUCCESS。

方法1:检查Scrapy运行状态主动抛出异常

在爬虫结束后,通过Scrapy的stats数据判断是否存在错误,若有则主动抛出异常触发Airflow任务失败:

def run_crawler(**kwargs):
    project_settings = set_project_settings({
        'FEEDS': {
            f'{kwargs["bucket"]}%(time)s.{kwargs["format"]}': {
                'format': kwargs["format"],
                'encoding': 'utf8',
                'store_empty': kwargs["store_empty"]
            }
        }
    })
    
    print("Project settings: ")
    pprint(project_settings.attributes.items())

    set_connection("airflow", kwargs["gcs_connection_id"])
    
    process = CrawlerProcess(project_settings)
    # 保存crawler实例用于后续获取stats
    crawler = process.crawl(spider.MySpider)

    print("Starting crawler...")
    # 使用run()替代start(),run()会阻塞直到爬虫完成
    process.run()
    
    # 获取爬虫运行统计数据
    stats = crawler.spider.stats.get_stats()
    
    # 检查关键错误指标:feed导出失败数、错误日志数
    feed_failed = stats.get('feedexport/failed_count/GCSFeedStorage', 0)
    error_count = stats.get('log_count/ERROR', 0)
    
    if feed_failed > 0 or error_count > 0:
        raise RuntimeError(f"爬虫运行出错:feed导出失败{feed_failed}次,错误日志{error_count}条")

方法2:捕获特定异常并重新抛出

针对明确的GCS凭证异常,直接捕获后重新抛出,确保Airflow感知:

from google.auth.exceptions import DefaultCredentialsError

def run_crawler(**kwargs):
    try:
        project_settings = set_project_settings({
            'FEEDS': {
                f'{kwargs["bucket"]}%(time)s.{kwargs["format"]}': {
                    'format': kwargs["format"],
                    'encoding': 'utf8',
                    'store_empty': kwargs["store_empty"]
                }
            }
        })
        
        print("Project settings: ")
        pprint(project_settings.attributes.items())

        set_connection("airflow", kwargs["gcs_connection_id"])
        
        process = CrawlerProcess(project_settings)
        crawler = process.crawl(spider.MySpider)

        print("Starting crawler...")
        process.run()
        
        # 额外检查stats确保无其他隐藏错误
        stats = crawler.spider.stats.get_stats()
        if stats.get('log_count/ERROR', 0) > 0:
            raise RuntimeError("爬虫运行过程中出现未捕获错误")
            
    except DefaultCredentialsError as e:
        # 捕获GCS凭证异常并重新抛出,触发Airflow任务失败
        raise RuntimeError(f"GCS凭证验证失败:{str(e)}") from e

关键说明

  • 替换process.start()为process.run():start()仅启动 reactor 不阻塞,run()会等待爬虫完成,方便后续检查状态。
  • 利用Scrapy stats:stats包含运行全量指标,可精准判断异常场景。
  • 主动抛出异常:Airflow的PythonOperator会捕获函数内的任何异常,自动将任务标记为FAILED。

内容的提问来源于stack exchange,提问作者no-stale-reads

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 04:42:54