恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
用100行Python写Mini-Airflow:实现爬虫任务依赖、重试与日志编排
首页
资讯中心
/
用100行Python写Mini-Airflow:实现爬虫任务依赖、重试与日志编排
用100行Python写Mini-Airflow:实现爬虫任务依赖、重试与日志编排
发布时间:2026/9/9 11:48:50
做爬虫项目任务一多你就会开始想一个事列表页抓完了才能抓详情页详情页入库之后才能做数据清洗这些环节怎么串起来自动跑我最早是写一个main.py从上到下硬编码顺序执行跑一会儿崩了从断点接着跑又得手动改代码。后来也试过部署一套完整的调度框架结果光维护调度器本身的时间就快赶上写业务代码了。最后我花了大半天用大概100行Python写了一个本地化的Mini-Airflow任务依赖、失败重试、执行日志全都有了。这篇文章就把这套方案完整拆开讲代码可以直接抄走改适合中小型爬虫项目、脚本自动化、数据采集管线的同学参考。1. 为什么自己写调度器而不是直接上Airflow1.1 调度器到底解决了爬虫里的什么问题爬虫项目看着简单实际落地时最耗时间的不是单个爬虫函数怎么写而是多个任务怎么编排。一个典型的采集流程通常是抓取列表页拿URL再去请求详情页解析字段最后把数据清洗后写入数据库或文件。这三步天然有先后顺序详情页的请求依赖列表页的结果清洗入库又依赖详情页的结果。如果你只有一个脚本手动维护顺序第一次跑没问题第二次多了增量任务、失败重试代码就开始变成一坨if else。调度器解决的正是这种任务编排问题。它把每个业务步骤抽象成一个独立任务任务之间声明依赖关系由调度器统一决定执行顺序。A没跑完B就不该启动B失败了要么重试要么标记失败而不是让下游无脑继续跑。这个思路在数据工程里叫DAG调度Airflow是其中的行业标准但爬虫项目往往不需要那么重的机制。1.2 工业级调度器的重和Mini方案从哪里切入Airflow专业、强大有Web界面、定时调度、分布式执行器、丰富的监控告警但对一个跑在本地的爬虫项目来说部署成本是真不低。你要装Python依赖、初始化数据库、配置Executor、启动Scheduler和Web Server一套下来至少多花半天时间。而且Airflow的任务实例状态存数据库和爬虫代码放一起反而让项目变臃肿。我需要的其实只有三个能力任务依赖、失败重试、日志记录。这三个能力本质上不需要任何外部依赖用Python原生的dict、set、logging就能实现。100行代码比引入一套框架更容易维护出了问题自己就能排查。这也是Mini-Airflow这类小工具存在的价值复杂工具解决复杂问题但普通人遇到的大部分问题其实用不上那套复杂度。1.3 什么时候该用这类Mini方案什么时候别硬凑Mini方案适合任务数量在几个到十几个之间、任务之间有明确的先后关系或并行关系不大、跑在单机上的场景。比如爬虫列表页→详情页→入库或者每日定时统计数据、生成报表再或者几个爬虫脚本之间互相补充数据。这种场景下你需要的只是在一个进程里把所有任务跑完跑完就走。如果你的任务数量到了几十个、需要跨机器调度、需要挂起和恢复历史任务状态、需要多人协作维护复杂的调度拓扑那还是老老实实用Airflow或Prefect吧。Mini方案没有持久化状态进程重启后一切从头来它定位是轻量、够用、贴近业务不是替代专业调度器。2. 设计思路任务依赖的本质是DAG2.1 核心抽象任务、依赖图、执行器实现一个小型调度器第一步是抽象出三个东西任务Task、依赖图Graph、执行器Executor。任务是业务函数本身带上任务ID、重试次数、重试间隔这些元信息。依赖图记录每个任务依赖谁底层用邻接表表示一个dictkey是任务IDvalue是这个任务依赖的任务ID集合。执行器负责按照依赖图算出执行顺序逐个运行任务运行失败就按策略重试。这和你爬山时的路线规划很像。多个景点之间有先后关系你得先到山口才能去山腰到了山腰才能去山顶。那地图本身就是一个图结构你的行走顺序就是对这个图做拓扑排序。2.2 拓扑排序怎么知道A得在B前面跑有向无环图DAG里拓扑排序能保证每个节点都在它的所有前置节点之后出现。对于依赖关系B依赖A排序结果里A一定在B前面。这样执行器只需要按排序后的列表顺序跑任务依赖天然满足。拓扑排序有两种经典算法Kahn算法和DFS。对于这个场景Kahn算法更好理解思路是先统计每个任务的入度入度为0的任务没有前置依赖可以立即执行执行完一个任务后把依赖它的所有任务的入度减1如果减到0就放进可执行队列循环直到所有任务都排进去。如果最后还有任务没排进去说明图里有环这就是检测到循环依赖的逻辑基础。2.3 功能取舍为什么少做反而好用写这个Mini调度器时我反复提醒自己一件事不要顺手把功能越加越多。网上有很多开源的迷你调度器写着写着就加了并发执行、定时触发、任务状态持久化、甚至Web UI代码量翻倍调试复杂度也跟着翻倍。一旦加了并发你就得处理线程安全、共享变量、任务间竞态这些坑比业务代码还难搞。我最终做的取舍是纯单线程顺序执行不加并发不做持久化进程结束状态就丢不做定时触发只做手动触发或配合系统级定时任务比如crontab不做Web界面日志就是最好的观测手段。这四个不做换来的是代码短、可读性强、出问题只需要看几十行就能定位。这种取舍在真实项目里非常重要功能边界清晰比功能多但没人敢改的代码强得多。2.4 任务之间怎么共享数据任务之间往往不是完全独立的。详情页解析需要列表页抓到的URL列表入库需要详情页解析出的字段。那任务之间怎么传递数据最简单的方案是让任务函数返回值由调度器把返回值存到一个共享的上下文字典里。下游任务执行时可以通过任务ID拿到上游任务的返回值。这就够了。不需要引入复杂的参数注入机制一个context字典打天下。具体实现是在执行器里维护一个results字典任务跑完后把返回值存进去执行下一个任务前先检查依赖任务是否都在results里如果不在就直接报错防止依赖没跑完就往下走的情况。3. 完整代码实现与逐段拆解3.1 基础设施日志与任务注册直接上代码先看整体结构。我把核心类命名为MiniAirflow全部代码放在一个文件里。import time import traceback import logging from collections import defaultdict, deque from typing import Callable, Dict, Any, Set class MiniAirflow: def __init__(self, namemini_airflow, log_fileNone): self.name name self.tasks: Dict[str, dict] {} self.dependencies: Dict[str, Set[str]] defaultdict(set) self.results: Dict[str, Any] {} self.logger self._setup_logger(log_file)这段做了三件事初始化任务字典、初始化依赖表、初始化日志器。tasks字典里存的是任务ID到任务配置的映射配置包括业务函数、重试次数、重试间隔。dependencies用defaultdict(set)来存依赖关系好处是直接add不需要先判断key是否存在。results是任务结果共享的context后面数据传递靠它。看日志器初始化我不建议直接用全局的logging而是给每个调度器实例建独立的logger这样日志可以写入不同文件互不干扰。def _setup_logger(self, log_file): logger logging.getLogger(self.name) logger.setLevel(logging.INFO) fmt logging.Formatter( %(asctime)s [%(levelname)s] %(message)s, datefmt%Y-%m-%d %H:%M:%S ) if log_file: fh logging.FileHandler(log_file, encodingutf-8) fh.setFormatter(fmt) logger.addHandler(fh) sh logging.StreamHandler() sh.setFormatter(fmt) logger.addHandler(sh) return logger这里有个细节如果重复运行脚本FileHandler会不断往同一个文件追加日志容易混乱。我实际使用时会按时间戳生成日志文件名比如crawler_20250101_100000.log每次运行一个文件排查问题很方便。3.2 任务装饰器与依赖注册为了让代码用起来像Airflow那样优雅我用装饰器注册任务。每个业务函数只要加上mini.task(...)就自动被登记到调度器里。def task(self, task_id: str, retries: int 3, retry_delay: int 5): def decorator(func: Callable): self.tasks[task_id] { func: func, retries: retries, retry_delay: retry_delay, } return func return decorator def set_dependency(self, task_id: str, depends_on: str): self.dependencies[task_id].add(depends_on)task这个装饰器接收三个参数task_id是任务唯一标识retries是最大重试次数retry_delay是失败后等待的秒数。这样业务代码里只要写清楚任务ID和重试策略就行执行细节全部交给调度器。set_dependency是声明依赖的核心方法含义是task_id这个任务要跑必须先跑完depends_on。这里我用depends_on这个参数名比parent或者upstream更直白。实际使用时一个任务可以依赖多个上游任务比如C依赖A和B那就要调用两次set_dependency。3.3 拓扑排序与循环依赖检测排序部分我用Kahn算法的简化版本。先统计每个任务的入度再按入度从0开始往出排。def _topological_sort(self) - list: in_degree {tid: 0 for tid in self.tasks} for tid, deps in self.dependencies.items(): for dep in deps: if dep not in self.tasks: raise KeyError(f依赖的任务 {dep} 不存在) in_degree[tid] 1 queue deque([tid for tid in self.tasks if in_degree[tid] 0]) order [] while queue: node queue.popleft() order.append(node) for tid, deps in self.dependencies.items(): if node in deps: in_degree[tid] - 1 if in_degree[tid] 0: queue.append(tid) if len(order) ! len(self.tasks): raise RuntimeError(检测到循环依赖任务无法正常排序) return order注意两个地方。第一依赖的任务必须存在如果写了个set_dependency(fetch_detail, fetch_list)但fetch_list这个任务ID并不存在这里会抛出KeyError帮你尽早发现配置错误。第二如果存在循环依赖比如A依赖B、B又依赖A那么两个节点的入度永远不可能变成0最后order的长度会小于tasks的总数。这一步的判断是整个调度器的安全性保障。为了让大家直观理解我加了一个简化版的依赖校验和拓扑排序测试。假设任务A、B、CC依赖A和B那in_degree初始为{A:0, B:0, C:2}队列先弹A或B执行完A后C的入度减为1再执行B后C入度减为0被加入队列最终顺序是[A, B, C]任意方向都满足依赖关系。3.4 执行引擎失败重试与结果收集核心的run方法就是把拓扑排序结果和任务执行串起来做失败重试、结果收集、日志记录。完整代码如下def run(self): self.logger.info(f调度器 {self.name} 开始运行) order self._topological_sort() self.logger.info(f任务执行顺序: { - .join(order)}) for tid in order: deps self.dependencies.get(tid, set()) missing [d for d in deps if d not in self.results] if missing: self.logger.error(f任务 {tid} 的依赖 {missing} 未成功执行跳过) raise RuntimeError(f依赖未满足: {tid} 缺少 {missing}) task_info self.tasks[tid] func task_info[func] retries task_info[retries] delay task_info[retry_delay] for attempt in range(1, retries 1): try: start time.time() self.logger.info(f开始执行任务 [{tid}] (第 {attempt}/{retries} 次尝试)) result func(self.results) elapsed time.time() - start self.results[tid] result self.logger.info(f任务 [{tid}] 执行成功, 耗时 {elapsed:.2f} 秒) break except Exception as e: self.logger.error(f任务 [{tid}] 执行失败: {e}) if attempt retries: self.logger.info(f任务 [{tid}] 将于 {delay} 秒后重试) time.sleep(delay) else: self.logger.error( f任务 [{tid}] 重试 {retries} 次后仍失败\n{traceback.format_exc()} ) raise self.logger.info(全部任务执行完毕) return self.results这里有两个设计细节值得展开。第一把self.results作为参数传给了任务函数下游任务就能通过self.results[上游任务ID]拿到上游结果。这个设计比Airflow的XCom简单但足够通用。第二重试策略是串行重试即失败后sleep几秒再跑对爬虫场景特别合适因为很多失败是网络抖动或目标站点临时限流隔几秒重试就能恢复。如果任务本身有内部状态也可以在上游为每个任务设置不同的重试参数。整个调度器的核心逻辑到这里就结束了从导入到头尾完全统计代码量大致在100行左右。对比Airflow动辄上万的代码量这个体量对爬虫项目非常友好。4. 实操过程跑通一套真实的爬虫调度4.1 设计任务链我先用一个电商商品采集的例子来演示。任务链设计为四步fetch_category获取分类页模拟请求分类列表返回分类链接列表。fetch_product_list根据分类链接抓商品列表页返回商品详情页URL列表。fetch_product_detail逐个请求详情页解析价格、标题、描述。save_to_database把解析结果整理后写入文件或数据库。这四个任务有明确的依赖关系fetch_product_list依赖fetch_categoryfetch_product_detail依赖fetch_product_listsave_to_database依赖fetch_product_detail。同时fetch_category没有上游可以最先执行save_to_database没有下游是最后一个执行。用set_dependency配置依赖时我会给每个任务配置不同的重试策略。网络请求类任务列表页、详情页重试次数多一些重试间隔也长一些因为网络抖动是偶发的数据库写入类任务重试次数少一些因为如果是数据本身有问题重试再多也没用。4.2 组装完整代码下面是把前面的调度器和具体爬虫任务拼接在一起的完整示例。实际项目里你只需要替换任务函数内的逻辑调度器部分完全不用改。import random import time from mini_airflow import MiniAirflow mini MiniAirflow(shop_crawler, log_fileshop_crawler.log) mini.task(fetch_category, retries3, retry_delay2) def fetch_category(ctx): 模拟抓取分类页返回分类链接列表 正常情况应该是 requests.get BeautifulSoup 解析 print(抓取分类页...) time.sleep(0.5) return [fhttps://shop.example.com/category/{i} for i in range(1, 4)] mini.task(fetch_product_list, retries3, retry_delay5) def fetch_product_list(ctx): categories ctx[fetch_category] print(f根据 {len(categories)} 个分类抓取商品列表...) product_urls [] for cat in categories: time.sleep(0.3) product_urls.extend([ f{cat}/product/{j} for j in range(1, 4) ]) return product_urls mini.task(fetch_product_detail, retries2, retry_delay10) def fetch_product_detail(ctx): product_urls ctx[fetch_product_list] print(f抓取 {len(product_urls)} 个商品详情...) results [] for url in product_urls: time.sleep(0.2) # 模拟随机失败一次验证重试机制 if random.random() 0.15: raise RuntimeError(f请求失败: {url}) results.append({ url: url, title: f商品-{url[-1]}, price: random.randint(10, 999) }) return results mini.task(save_to_database, retries1, retry_delay3) def save_to_database(ctx): products ctx[fetch_product_detail] print(f保存 {len(products)} 条商品数据...) with open(products.txt, w, encodingutf-8) as f: for p in products: f.write(f{p[title]} | {p[price]} | {p[url]}\n) return fsaved {len(products)} items mini.set_dependency(fetch_product_list, fetch_category) mini.set_dependency(fetch_product_detail, fetch_product_list) mini.set_dependency(save_to_database, fetch_product_detail) mini.run()任务函数统一接收ctx这个参数也就是调度器内部维护的results字典。每个任务通过ctx[上游任务ID]拿上游数据。这样代码结构清晰任务之间没有任何隐式共享变量符合可测试、可维护的要求。这里random.random() 0.15模拟了详情页抓取时的随机失败用来验证重试机制的效果。配合日志你能清楚看到哪个任务在哪个时间点失败了多少次、第几次成功。4.3 运行结果与日志解读运行上面的代码控制台会输出类似下面的日志2025-01-01 10:00:01 [INFO] 调度器 shop_crawler 开始运行 2025-01-01 10:00:01 [INFO] 任务执行顺序: fetch_category - fetch_product_list - fetch_product_detail - save_to_database 2025-01-01 10:00:01 [INFO] 开始执行任务 [fetch_category] (第 1/3 次尝试) 2025-01-01 10:00:01 [INFO] 任务 [fetch_category] 执行成功, 耗时 0.51 秒 2025-01-01 10:00:02 [INFO] 开始执行任务 [fetch_product_list] (第 1/3 次尝试) ... 2025-01-01 10:00:05 [ERROR] 任务 [fetch_product_detail] 执行失败: 请求失败: https://shop.example.com/category/2/product/2 2025-01-01 10:00:06 [INFO] 任务 [fetch_product_detail] 将于 10 秒后重试 2025-01-01 10:00:16 [INFO] 开始执行任务 [fetch_product_detail] (第 2/2 次尝试) 2025-01-01 10:00:19 [INFO] 任务 [fetch_product_detail] 执行成功, 耗时 2.95 秒 ... 2025-01-01 10:00:21 [INFO] 全部任务执行完毕日志的关键价值在于可复盘。可以看出每一次重试的开始时间、耗时、失败原因甚至可以按任务ID过滤日志来追踪单个任务的多次尝试。如果你用log_file参数指定了日志文件那么控制台和文件都会有同步输出。这里我建议把日志级别设为INFO任务里的print只能到控制台不会进日志文件。所以业务函数里如果有重要的输出可以用logging.getLogger(shop_crawler).info(...)或者统一把print改成ctx.logger之类的自定义输出保证日志文件内容完整。4.4 扩展任务结果持久化与断点恢复Mini方案虽然不持久化任务状态但有一个非常实用的扩展把self.results定期用json或pickle写进本地文件。如果任务链很长跑了一半进程崩溃了下次启动时可以先加载上次的结果跳过已经成功的任务。这个扩展不需要改调度器核心只要在run方法里加几行import json, os def _save_checkpoint(self, tid): tmp {k: v for k, v in self.results.items() if isinstance(v, (str, int, list, dict))} with open(f{self.name}_checkpoint.json, w, encodingutf-8) as f: json.dump(tmp, f, ensure_asciiFalse, indent2) def _load_checkpoint(self): path f{self.name}_checkpoint.json if os.path.exists(path): with open(path, r, encodingutf-8) as f: self.results.update(json.load(f))我实际使用中遇到网络爬虫中途断开断点恢复能省掉重爬大量页面。不过要提醒一下results里的数据如果是自定义对象json序列化不了这种情况需要用pickle或者只持久化任务状态、不持久化数据本身。5. 常见问题与排查技巧实录5.1 循环依赖为什么一言不合就出现循环依赖是写调度配置时最容易犯的错尤其是任务多了以后。比如你写的时候想着清洗依赖入库入库也依赖清洗两个任务互相等待拓扑排序就会报错。我的排查经验是先打印出所有set_dependency调用沿着依赖链画一遍只要发现从某个任务出发能绕一圈回来那就是环。实际操作中我会在set_dependency里加一条校验不允许任务依赖自己这是最常见的隐性环。如果配置了A - B - C - A这种链条我的_topological_sort会抛出RuntimeError并把参与任务总数和已排序任务数都打出来方便定位。5.2 重试参数到底怎么设置才合理很多教程告诉你失败重试3次但实际场景里重试策略要分任务来定。网络请求类的任务重试间隔长一点5到30秒给目标站点留出恢复时间数据库写入类的任务间隔短一点2到5秒重试次数少一点因为大概率是数据问题。我自己常用的策略是请求类retries3, retry_delay10解析类retries2, retry_delay3入库类retries1, retry_delay2。重试不是越多越好。如果目标站点已经封了你的IP重试10次也是白搭。所以我的调度器没有设计指数退避而是固定间隔重试因为爬虫场景下简单的策略往往最可控。如果你想做指数退避在run方法里把sleep(delay)改成sleep(delay * attempt)就行加了指数退避后重试效果会更有节奏感。5.3 PyCharm里只显示Process finished with exit code 0没输出内容怎么办这个现象太常见了。你点运行控制台只显示一行Process finished with exit code 0中间啥内容都没有。出现这个大部分情况是程序确实跑完了但输出被吞了或者根本没执行到你预期的地方。在调度器这个场景里常见原因有三个一是日志级别设置太高比如logger.setLevel(logging.ERROR)INFO级别的日志自然不显示二是任务函数内部用了print但调度器代码里加了日志handler覆盖了输出流导致print的内容没刷出来三是你有多个logger实例日志写到了文件而不是控制台。我的排查顺序是先看log_file指定的日志文件里有没有内容有内容就说明代码跑了只是控制台没显示然后在任务函数第一行加一个print(task start)看是否执行到最后检查logger配置确认StreamHandler确实被添加了。如果还不行把logger.setLevel(logging.DEBUG)看看是不是有早期异常被吞了。5.4 日志重复输出和中文乱码问题日志重复输出主要原因是你重复执行脚本时前一次运行的FileHandler没有释放。用PyCharm反复运行同一个logger名字会挂多个handler导致日志打两遍。解决办法是在_setup_logger里先清空已有handlerlogger.handlers.clear()或者用logger.removeHandler(handler)精确清理。这个细节别看小实际调试时能省不少时间。中文乱码问题几乎都是文件编码没指定。我在初始化logging.FileHandler时明确写了encodingutf-8但Windows上默认编码是gbk如果没指定encoding中文日志就会变成乱码。这个问题排查起来很隐蔽因为错误信息里根本不会提示编码问题。如果你在Windows上遇到乱码优先检查FileHandler的编码参数其次检查日志文件是用什么编码打开的。另外提醒一下在爬虫任务里如果要把数据写入文件也一定要带encodingutf-8否则Windows下写入中文会报UnicodeEncodeError。这也是新手写爬虫最容易遇到的问题。5.5 任务函数改了日志却显示旧代码在跑这个坑我踩过好几次。Python的模块缓存机制会让旧的.pyc文件被重新使用尤其是PyCharm里热更新不彻底的时候。遇到修改代码后日志输出还是旧逻辑先重启解释器再清一遍__pycache__目录。在命令行跑一下find . -name *.pyc -delete问题基本就解决了。如果你已经改成了重启大法还是不行那就要检查是不是有多个Python环境跑的时候用的不是当前项目解释器。6. 后续扩展让Mini调度器更贴近实际项目我这个Mini方案在项目里跑了一段时间之后稳定性和维护成本都非常理想。在此基础上你还可根据实际需求做几个小扩展。加定时触发在main.py里用while True time.sleep循环判断当前时间是否到达目标时间点到了就调用mini.run()。或者直接配系统crontab更省事。加并发执行把所有独立的入度为0的任务丢到线程池里并行执行。要处理共享变量和异常收集代码量会增加但如果你有几组互不依赖的爬虫任务提速效果很明显。加Webhook通知任务失败时用requests.post把失败信息发到企业微信、钉钉或邮件。接入很轻量对长时间运行的脚本尤其有用。我个人在实际项目里目前这套Mini-Airflow承担了三个爬虫项目的数据采集任务。相比之前每个项目各自硬编码执行顺序统一调度之后新增任务只要加一个函数注册依赖几乎不用动其他代码。任务失败了看日志就能定位日志文件按时分秒命名归档和排查都很方便。如果你也需要编排多个爬虫任务不妨从这套100行的方案开始它简单到让你可以把精力集中在业务逻辑本身而不是调度器上的各种新坑。