编程式调度作业防重复:Spark定时任务因YARN资源竞争重复触发解决方案咨询
碰到过一模一样的问题!之前用Cron调度Spark作业存物联网传感器数据到Hive,也是因为YARN资源紧张导致作业延迟,Cron还一个劲触发重复任务,把集群搞得更卡。分享几个我在生产环境验证过的方案:
核心思路是让同一时间间隔的作业,同一时刻只能有一个实例在运行,用分布式锁来做互斥控制,避免Cron重复触发。推荐用ZooKeeper(Curator客户端)或者Redis实现,这里给个ZooKeeper的Scala伪代码示例:
import org.apache.curator.framework.CuratorFrameworkFactory import org.apache.curator.retry.ExponentialBackoffRetry import org.apache.curator.framework.recipes.locks.InterProcessMutex import java.util.concurrent.TimeUnit object SensorJobLock { def main(args: Array[String]): Unit = { val zkServers = "zk-node1:2181,zk-node2:2181" val jobInterval = args(0) // 比如传入"30min",作为锁的标识 val lockPath = s"/locks/spark_sensor_job_${jobInterval}" // 初始化Curator客户端 val curator = CuratorFrameworkFactory.newClient(zkServers, new ExponentialBackoffRetry(1000, 3)) curator.start() val lock = new InterProcessMutex(curator, lockPath) try { // 尝试获取锁,10秒内获取不到就退出 if (lock.acquire(10, TimeUnit.SECONDS)) { println("成功获取锁,开始执行Spark作业...") // 初始化SparkSession并执行业务逻辑 val spark = org.apache.spark.sql.SparkSession.builder() .appName(s"SensorDataIngest_${jobInterval}") .enableHiveSupport() .getOrCreate() // 这里写你的传感器数据读取、清洗、写入Hive表的逻辑 // spark.read.format("...").load().write.mode("append").saveAsTable("hive_db.sensor_table") spark.stop() println("作业执行完成,释放锁") } else { println("检测到同类型作业正在运行,当前实例退出") System.exit(0) } } finally { // 确保锁被释放,避免死锁 if (lock.isAcquiredInThisProcess) { lock.release() } curator.close() } } }
这个方案的优势是跨节点可靠,不管你的作业提交到哪个节点,都能通过分布式锁判断是否有同类作业在跑,不会出现漏判的情况。
如果你的集群有多个调度任务需要管理,直接换掉Cron用Airflow这类工作流工具更省心。Airflow天然支持任务的状态依赖,可以配置成「上一个周期的作业完成后,再触发下一个」,从根源上避免重复触发。
给个Airflow的DAG配置示例(Python):
from airflow import DAG from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from datetime import datetime, timedelta # 默认参数配置 default_args = { 'owner': 'data_team', 'depends_on_past': True, # 强制依赖上一次任务的执行状态 'wait_for_downstream': True, # 确保下游任务(如果有的话)完成再触发下一轮 'start_date': datetime(2024, 5, 1), 'retries': 0, # 因为要避免重复,这里不配置重试,按需调整 'retry_delay': timedelta(minutes=5), } # 定义DAG,30分钟间隔执行 dag = DAG( 'sensor_data_to_hive', default_args=default_args, description='定时将传感器数据写入Hive表', schedule_interval='*/30 * * * *', catchup=False, # 禁止补跑历史任务,避免集群积压 max_active_runs=1, # 同一时间只允许一个DAG实例运行 ) # Spark提交任务 spark_ingest_task = SparkSubmitOperator( task_id='spark_sensor_ingest_30min', application='/data/jars/sensor_ingest_job.jar', conn_id='yarn_cluster', # 提前在Airflow配置好YARN连接 executor_cores=2, executor_memory='4g', num_executors=3, dag=dag, ) spark_ingest_task
Airflow会自动监控Spark作业在YARN上的状态,上一个作业没完成的话,下一个会处于等待状态,绝对不会重复提交。而且还能直观看到任务的运行历史、失败日志,比Cron好用太多。
如果不想引入额外组件,最简单的方式是在提交作业前,先检查YARN上是否有同名的作业在运行。可以写个shell脚本,让Cron触发这个脚本,而不是直接触发spark-submit:
#!/bin/bash # 定义作业名称,要和Spark作业的appName一致 JOB_NAME="SensorDataIngest_30min" # 检查YARN上是否有处于RUNNING或ACCEPTED状态的同名作业 RUNNING_JOBS=$(yarn application -list 2>/dev/null | grep "$JOB_NAME" | grep -E "(RUNNING|ACCEPTED)") if [ -z "$RUNNING_JOBS" ]; then echo "未检测到运行中的作业,开始提交..." spark-submit \ --class com.yourcompany.SensorIngest \ --master yarn \ --deploy-mode cluster \ --executor-cores 2 \ --executor-memory 4g \ /data/jars/sensor_ingest_job.jar else echo "作业 $JOB_NAME 正在运行/等待资源,本次不提交" exit 0 fi
这个方案的优点是轻量无依赖,只需要执行脚本的用户有YARN的查看权限就行。缺点是如果作业名称重复可能会误判,所以要确保每个时间间隔的作业用唯一的appName。
内容的提问来源于stack exchange,提问作者Ashwini

