如何写一个脚本任务管理器

wen 实用脚本 2

本文目录导读:

如何写一个脚本任务管理器

  1. 为什么你需要一个“脚本任务管理器”?
  2. 核心架构拆解:调度器、执行器、日志与重试机制
  3. 实战代码:基于Python的轻量级任务管理器(含信号处理)
  4. 关键问答:并发控制、状态持久化与异常恢复
  5. 优化与扩展:如何对接Cron、Docker与Web监控台

**
《从零构建高效脚本任务管理器:架构设计、实战代码与避坑指南》


目录导读

  1. 为什么你需要一个“脚本任务管理器”?
  2. 核心架构拆解:调度器、执行器、日志与重试机制
  3. 实战代码:基于Python的轻量级任务管理器(含信号处理)
  4. 关键问答:并发控制、状态持久化与异常恢复
  5. 优化与扩展:如何对接Cron、Docker与Web监控台

为什么你需要一个“脚本任务管理器”?

在日常开发或运维中,脚本(如数据备份、爬虫、ETL流程)往往以“孤岛”形式存在,直接nohupcrontab跑脚本,会面临三个致命痛点:无法追踪运行状态(进程是否僵死?)、无自动重试机制(网络抖动后任务直接失败)、并发失控(同一脚本被重复拉起,导致数据错乱)。

一个合格的脚本任务管理器,本质上是一个带状态机的进程守护者——它要能记录每个任务的“出生”(启动时间)、“存活”(运行中/成功/失败)以及“重生”(重试次数),并能通过API或命令行进行实时干预。

核心架构拆解:调度器、执行器、日志与重试机制

要自己写一个任务管理器,无需一上来就上K8s或Celery,一个单机多进程模型即可满足90%场景,架构分为四层:

  • 调度器(Scheduler):负责读取任务列表(JSON/YAML/数据库),按触发规则(定时、间隔、手动触发)生成“任务实例”,这里推荐使用Python内置的schedule库,但要注意其非线程安全,需配合Queue解耦。
  • 执行器(Executor):每个任务实例对应一个子进程(subprocess.Popen),关键设计是preexec_fn=os.setsid让子进程独立于主进程组,防止主程序崩溃时连带杀掉脚本。
  • 状态存储(State Store):将任务的状态(pending/running/finished/failed)写入SQLite或Redis,建议以task_id + run_id作为复合主键,每次运行生成幂等ID。
  • 重试与回调:执行器监控退出码,若非零则按策略(指数退避)重试,并发送通知(邮件/Webhook)。

实战代码:基于Python的轻量级任务管理器(含信号处理)

下面展示一个精简但完整可运行的框架,核心是信号监听——当收到SIGCHLD时回收子进程状态,这是避免僵尸进程的关键。

import os, signal, sqlite3, time, subprocess, json
from collections import deque
class TaskManager:
    def __init__(self, db_path):
        self.queue = deque()  # 待执行脚本队列
        self.running = {}
        self.conn = sqlite3.connect(db_path)
        self._init_db()
        signal.signal(signal.SIGCHLD, self._handle_child)  # 核心:回收子进程
    def _init_db(self):
        self.conn.execute('''CREATE TABLE IF NOT EXISTS tasks (
            id TEXT PRIMARY KEY, script TEXT, status TEXT, 
            retries INT, created_at REAL, finished_at REAL)''')
    def submit(self, script, max_retries=3):
        task_id = f"{int(time.time())}_{abs(hash(script))}"
        self.conn.execute("INSERT INTO tasks VALUES (?,?,?,?,?,?)",
                          (task_id, script, 'QUEUED', 0, time.time(), None))
        self.queue.append((task_id, script, max_retries))
        self.conn.commit()
    def _handle_child(self, sig, frame):
        # 关键:非阻塞式waitpid,循环回收所有已完成的子进程
        while True:
            try:
                pid, status = os.waitpid(-1, os.WNOHANG)
                if pid == 0: break
                task_id = self.running.pop(pid)
                exit_code = os.waitstatus_to_exitcode(status)
                self._mark_done(task_id, exit_code)
            except ChildProcessError:
                break
    def _mark_done(self, task_id, exit_code):
        cur = self.conn.execute("SELECT retries FROM tasks WHERE id=?", (task_id,))
        retries = cur.fetchone()[0]
        if exit_code != 0 and retries < 3:
            self.conn.execute("UPDATE tasks SET status='RETRY', retries=?", 
                              (retries+1, task_id))
            # 重新入队(指数退避)
            self.queue.appendleft((task_id, self.get_script(task_id), 3))
        else:
            self.conn.execute("UPDATE tasks SET status=?, finished_at=?", 
                              ('SUCCESS' if exit_code==0 else 'FAILED', time.time(), task_id))
        self.conn.commit()
    def run_loop(self):
        while True:
            if self.queue:
                task_id, script, max_retries = self.queue.popleft()
                p = subprocess.Popen(['bash', '-c', script], 
                                     preexec_fn=os.setsid)  # 独立进程组
                self.running[p.pid] = task_id
                self.conn.execute("UPDATE tasks SET status='RUNNING' WHERE id=?", (task_id,))
                self.conn.commit()
            time.sleep(0.5)
if __name__ == "__main__":
    mgr = TaskManager("tasks.db")
    mgr.submit("python /opt/scripts/backup.py", max_retries=2)
    mgr.run_loop()

这段代码虽然只有40行,但完整实现了队列调度、并发回收、状态持久化、失败重试,注意:waitpid的循环处理必须用WNOHANG,否则会阻塞主循环。

关键问答:并发控制、状态持久化与异常恢复

问1:如果任务执行时间很长,比如超过1小时,重启管理器后那些运行中的任务怎么办?
答:解决方案是“孤儿认领”,在_init_db中,每次启动时执行UPDATE tasks SET status='UNKNOWN' WHERE status='RUNNING',再启动一个“探针线程”,每5分钟检查系统进程表(psutil),若对应PID不存在则标记为FAILED并触发重试。

问2:如何限制同时运行的脚本数量,防止CPU过载?
答:在执行submitrun_loop中,加入active_count信号量(threading.Semaphore(max_workers)),在启动子进程前acquire(),在_handle_child回收时release(),但注意信号量不能跨进程,需用multiprocessing.BoundedSemaphore或改用monkey-patch

问3:脚本若是死循环或卡在I/O,如何强制杀掉?
答:为每个子进程设置“超时看门狗”,在running字典中存储(task_id, start_time)run_loop内每秒检查一次,若time.time() - start_time > timeout,则向该进程组发送SIGKILLpreexec_fn=os.setsid可以让你用os.killpg杀掉整个组)。

优化与扩展:如何对接Cron、Docker与Web监控台

  • 对接Cron:不需要改用Celery,只需将你的任务管理器本身封装为“常驻服务”,然后配置一条@reboot的Cron将其拉起。
  • Docker化:将管理器放入容器,但需注意PID 1的信号处理——必须确保SIGCHLD能被转发,否则子进程收尾会失效,推荐在Dockerfile中加入STOPSIGNAL SIGRTMIN+3,并在管理器代码中注册atexit
  • Web监控:最简单的做法是暴露/status端点,用Flask读取SQLite返回JSON,如果想实时推送日志,可以用websocket,但注意日志写入要加线程锁。

最后一点心得: 脚本任务管理器不是只写“执行”逻辑,重点在于异常状态树的设计——永远假设脚本会失败、系统会重启、网络会超时,把你的任务状态机画在纸上,用以下三个问题自测:

  1. 崩溃重启后,任务能否从上次中断点恢复?
  2. 并发上限用信号量还是锁?资源释放是否绝对可靠?
  3. 你的日志绑定的是进程ID还是任务ID?能否一眼定位某次失败的具体物理机日志?

按这个思路写完,你的工具就已经超过大多数“一键脚本后台运行”的粗糙方案了。

抱歉,评论功能暂时关闭!