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.
相关产品推荐
相关产品推荐

