如何设计Python Worker+Flask API避免API因数据采集阻塞?
嘿,我完全懂你这种靠试错摸爬滚打的感觉——没有专业背景搞架构确实容易卡壳在这些异步/并发的点上。咱们一步步来拆解你遇到的问题,给你一套能落地的方案,还有可参考的思路扩展到大规模场景。
核心思路:彻底解耦API与Worker,用队列做“中间缓冲”
你之前想到的线程和队列方向完全对,但问题出在没把API的请求处理和Worker的采集任务彻底分开。API的核心职责是快速响应外部请求,采集这种耗时的活儿,必须丢给后台独立的Worker去做,中间用消息队列来传递任务指令——这样API永远不会被采集任务拖垮。
第一步:选个对新手友好的消息队列工具
不用一开始就啃复杂的RabbitMQ,推荐先从**Redis Queue(RQ)**入手,它封装得极其简洁,几乎不用写底层队列逻辑,跟着示例就能跑起来。
第二步:最小可行实现(单API+单Worker)
我给你写一套能直接跑的代码结构,你可以对着改:
项目结构
your_project/ ├── app.py # Flask API服务 ├── worker.py # 后台采集Worker ├── tasks.py # 单独抽离的采集任务逻辑 └── requirements.txt
先装依赖
requirements.txt里放这些:
flask rq redis requests # 假设你用requests做网页采集,换成你用的工具就行
然后运行 pip install -r requirements.txt 安装。
1. 写采集任务逻辑(tasks.py)
把采集数据的代码单独抽出来,做成可被队列调度的函数:
import requests def fetch_web_data(target_url): # 这里替换成你的实际采集逻辑:请求网页、解析数据、存数据库等 try: resp = requests.get(target_url, timeout=15) resp.raise_for_status() # 举个例子:返回结构化的采集结果,或者直接写入数据库 result = { "url": target_url, "content_length": len(resp.text), "status": "success" } # 这里可以加数据库写入代码,比如用SQLAlchemy操作MySQL/PostgreSQL return result except Exception as e: return { "url": target_url, "error": str(e), "status": "failed" }
2. 写Flask API(app.py)
API只做一件事:接收请求,把任务丢进队列,立刻返回响应,绝对不自己执行采集:
from flask import Flask, request, jsonify from rq import Queue from redis import Redis import tasks app = Flask(__name__) # 连接本地Redis,初始化任务队列 redis_conn = Redis(host="localhost", port=6379) task_queue = Queue(connection=redis_conn) @app.route("/submit-collect", methods=["POST"]) def submit_collect_task(): req_data = request.get_json() if not req_data or "url" not in req_data: return jsonify({"error": "必须提供url参数"}), 400 # 把采集任务丢进队列,拿到任务ID返回给客户端 task = task_queue.enqueue(tasks.fetch_web_data, req_data["url"]) return jsonify({ "task_id": task.get_id(), "status": "已加入队列", "msg": "采集任务已提交,后台正在处理" }), 202 # 可选:加一个查询任务状态的接口,方便客户端跟踪结果 @app.route("/task-status/<task_id>", methods=["GET"]) def check_task_status(task_id): task = task_queue.fetch_job(task_id) if not task: return jsonify({"error": "任务不存在"}), 404 resp_data = { "task_id": task_id, "status": task.get_status(), "result": task.result if task.is_finished else None } return jsonify(resp_data) if __name__ == "__main__": # Flask默认就是多线程模式,生产环境别用debug模式 app.run(host="0.0.0.0", port=5000, threaded=True)
3. 启动Worker(worker.py)
单独启动一个进程,专门处理队列里的采集任务:
from rq import Worker, Queue, Connection from redis import Redis redis_conn = Redis(host="localhost", port=6379) if __name__ == "__main__": with Connection(redis_conn): # 监听默认队列,有多个队列的话可以在这里添加 worker = Worker([Queue()]) worker.work()
怎么运行这套代码?
- 先启动Redis服务(本地装了的话直接跑
redis-server) - 启动Worker:
python worker.py(这个进程会一直跑,处理队列任务) - 启动Flask API:
python app.py
这样一来,API每次接请求都是瞬间响应,采集的耗时工作全在Worker进程里跑,和API的线程池完全隔离——再也不会出现API被采集任务卡爆的情况。
第三步:扩展到多Worker+多Flask API
如果要放大规模,逻辑其实很简单,就是横向复制:
- 多Worker:直接启动多个
worker.py进程就行,RQ会自动把队列里的任务分配给空闲的Worker。服务器上可以用supervisor或者systemd来管理多个Worker进程,防止意外退出。 - 多Flask API:比如用Gunicorn启动多个Flask进程,只要所有API实例都连接同一个Redis队列,就能统一往队列里丢任务,Worker们一起处理,完全不冲突。
- 如果需要更复杂的功能(比如任务优先级、延迟执行、定时任务),可以换成
Celery+RabbitMQ,但Celery学习曲线比RQ陡一点,建议先把RQ玩熟了再升级。
可参考的开源项目思路
很多轻量级爬虫调度系统都用了这种架构:
- RQ官方的示例项目:里面有各种异步任务的典型用法,非常适合新手参考
- 一些开源的博客/内容管理系统:比如批量导入文章的功能,就是API接请求,后台Worker跑导入逻辑,你可以找这类项目的异步任务模块看代码
内容的提问来源于stack exchange,提问作者entalpia
相关产品推荐
相关产品推荐

