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

SparkSubmitOperator运行时触发'dag_id could not be found'错误及工作流停滞问题

排查Airflow中SparkSubmitOperator任务停滞的思路

从你给出的日志来看,任务只走到了Airflow提交Spark任务的步骤,没有后续的执行日志或者错误输出,这种情况大概率是Spark任务本身卡在了执行阶段,而非Airflow调度层面的问题。我之前处理过类似的场景,给你几个具体的排查方向:

  • 优先获取Spark任务的完整执行日志:Airflow的日志仅记录了它提交Spark任务的动作,实际Spark任务的运行细节(比如driver报错、executor状态异常)需要到Spark集群的日志系统中查找:

    • 如果是YARN集群,登录YARN ResourceManager的UI,找到对应Application ID的任务,查看driver和executor的日志;
    • 如果是Spark Standalone模式,进入Spark Master的UI定位任务,再查看对应的本地日志文件。
      这一步是核心,任务停滞的具体原因(比如S3拉取超时、权限不足、内存溢出)都会在Spark的专属日志里体现。
  • 验证S3访问权限:日志显示任务是用hdfs用户执行的,要确认这个用户是否具备目标S3桶的访问权限:

    • 如果是AWS EMR集群,检查EMR实例绑定的IAM角色是否包含S3的读权限;
    • 如果是手动配置AWS密钥,确认hdfs用户的环境变量或Spark配置中的fs.s3a.access.key和fs.s3a.secret.key是否正确;
    • 同时检查S3桶的Bucket Policy是否限制了访问的IP地址或用户范围。
  • 检查Spark任务的资源与配置:

    • 确认Spark任务的driver和executor内存、CPU分配是否充足,拉取大量S3数据时很容易因OOM导致任务卡住;
    • 针对S3数据拉取,建议配置S3A的优化参数,比如增大连接数fs.s3a.connection.maximum、调整预读范围fs.s3a.readahead.range,这些参数能提升S3数据读取效率,避免超时停滞。
  • 脱离Airflow单独测试Spark任务:直接用hdfs用户在服务器上执行对应的Spark Submit命令,比如:

    sudo -H -u hdfs spark-submit --class com.your.package.ImportCrawlJob --master yarn-cluster your-spark-job.jar
    

    如果单独运行也卡住,说明问题出在Spark任务本身,和Airflow无关;如果单独运行正常,再去排查Airflow的SparkSubmitOperator配置(比如提交模式、参数传递是否正确)。

附你提供的任务日志:

[2018-03-22 13:37:02,762] {models.py:1428} INFO - Executing <Task(SparkSubmitOperator): ImportCrawlJob> on 2018-03-22 13:37:00
[2018-03-22 13:37:02,763] {base_task_runner.py:115} INFO - Running: ['bash', '-c', 'sudo -H -u hdfs airflow run dag_extract_jobs Impo...

内容的提问来源于stack exchange,提问作者Aaron Z.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:18:19