You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何设计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()

怎么运行这套代码?

  1. 先启动Redis服务(本地装了的话直接跑 redis-server)
  2. 启动Worker:python worker.py(这个进程会一直跑,处理队列任务)
  3. 启动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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:05:08