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

基于Celery与Flask实现PDF分页拆分异步工作流的技术问询

解决方案:Celery+Flask 动态PDF处理工作流实现

核心逻辑

要实现「异步统计PDF页数→动态生成并行拆页任务→生成最终报告」的非阻塞工作流,关键是利用Celery的任务签名(signature)和组任务(group)、**和弦任务(chord)**特性,让count_pages的结果动态驱动后续任务的生成与执行。

具体实现

1. 定义基础任务

先实现三个核心任务,确保count_pages返回后续任务所需的关键参数:

from celery import Celery, group, chord
import pdfplumber
from PyPDF2 import PdfReader, PdfWriter
import os

# 初始化Celery实例
app = Celery('pdf_processor', broker='redis://localhost:6379/0')

@app.task
def count_pages(pdf_path):
    # 统计PDF页数逻辑
    with pdfplumber.open(pdf_path) as pdf:
        page_count = len(pdf.pages)
    # 返回页数+文件路径,供后续任务使用
    return {"page_count": page_count, "pdf_path": pdf_path}

@app.task(autoretry_for=(Exception,), retry_backoff=3, retry_kwargs={"max_retries": 2})
def split_page(pdf_path, page_num):
    # 拆分单页为独立PDF逻辑
    reader = PdfReader(pdf_path)
    writer = PdfWriter()
    writer.add_page(reader.pages[page_num - 1])  # 页码从1开始计数
    output_path = f"{os.path.splitext(pdf_path)[0]}_page_{page_num}.pdf"
    with open(output_path, 'wb') as out_file:
        writer.write(out_file)
    return output_path

@app.task
def generate_report(split_results):
    # 生成处理报告逻辑
    report_content = f"PDF处理完成:共拆分{len(split_results)}页\n拆分文件路径:\n" + "\n".join(split_results)
    report_path = f"{os.path.splitext(split_results[0])[0]}_report.txt"
    with open(report_path, 'w') as f:
        f.write(report_content)
    return report_path

2. 编排动态工作流

通过任务回调或链式任务,实现count_pages完成后自动触发并行拆页任务,最后执行报告生成:

方式一:回调触发和弦任务

在Flask上传接口中异步启动count_pages,通过link参数指定后续任务,动态生成拆页任务组:

from flask import Flask, request

flask_app = Flask(__name__)
UPLOAD_FOLDER = "/tmp/pdf_uploads"
os.makedirs(UPLOAD_FOLDER, exist_ok=True)

@flask_app.route('/upload-pdf', methods=['POST'])
def upload_pdf():
    if 'pdf' not in request.files:
        return "未上传文件", 400
    pdf_file = request.files['pdf']
    pdf_path = os.path.join(UPLOAD_FOLDER, pdf_file.filename)
    pdf_file.save(pdf_path)
    
    # 异步启动统计任务,不阻塞接口
    count_pages.apply_async(args=[pdf_path], link=trigger_split_workflow.s())
    return "PDF处理已启动", 202

@app.task
def trigger_split_workflow(count_result):
    page_count = count_result["page_count"]
    pdf_path = count_result["pdf_path"]
    # 生成对应页数的拆页任务组
    split_task_group = group(
        split_page.s(pdf_path, page_num) for page_num in range(1, page_count + 1)
    )
    # 用chord确保所有拆页任务完成后执行报告生成
    chord(split_task_group)(generate_report.s())

方式二:链式任务直接编排

直接在上传接口中定义链式任务,利用count_pages的结果动态生成任务组:

@flask_app.route('/upload-pdf', methods=['POST'])
def upload_pdf():
    # 省略文件保存逻辑...
    
    # 定义链式工作流:统计页数→生成拆页组→生成报告
    workflow = chain(
        count_pages.s(pdf_path),
        # 用lambda接收统计结果,动态生成拆页任务组
        (lambda res: group(
            split_page.s(res["pdf_path"], num) for num in range(1, res["page_count"] + 1)
        )) | generate_report.s()
    )
    # 异步执行整个工作流
    workflow.apply_async()
    return "PDF处理已启动", 202

3. 关键注意事项

  • 参数传递:确保count_pages返回的结果包含后续任务所需的所有信息(如文件路径、页数),避免依赖全局变量或不可靠的外部存储。
  • 异常处理:给拆页任务添加自动重试机制,应对临时的文件IO或PDF解析异常。
  • Broker性能:生产环境建议使用Redis或RabbitMQ作为Celery Broker,避免使用默认的SQLAlchemy(性能较差)。

内容的提问来源于stack exchange,提问作者maxime

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 12:06:12