基于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
相关产品推荐
相关产品推荐

