Flask多进程RSS/Atom解析报错:无法pickle本地对象
这个错误Can't pickle local object 'SQLAlchemy.init_app.<locals>.shutdown_session'其实很好理解:当你在多进程中传递Flask app或者关联了SQLAlchemy的对象时,pickle(Python多进程用来传递数据的序列化机制)尝试序列化shutdown_session这个函数——但它是SQLAlchemy.init_app内部定义的局部函数,pickle没法处理这类没有全局引用的对象。
下面给你几个可行的解决方案,按推荐程度排序:
方案1:让每个子进程独立初始化App和DB(最推荐)
Flask和SQLAlchemy都不是为跨进程共享设计的,每个进程应该有自己的App实例和数据库会话。我们可以重构代码,让子进程自己创建App上下文,只传递必要的数据(比如Feed的URL和ID,而不是整个模型实例或App对象)。
第一步:重构App初始化逻辑
先把App的创建抽成可复用的函数(通常放在app/__init__.py里):
from flask import Flask from flask_sqlalchemy import SQLAlchemy db = SQLAlchemy() def create_app(): app = Flask(__name__) # 加载配置,比如数据库连接信息 app.config['SQLALCHEMY_DATABASE_URI'] = 'your-db-uri' app.config['SQLALCHEMY_TRACK_MODIFICATIONS'] = False db.init_app(app) return app
第二步:修改并行解析函数
让函数在子进程内部创建App并查询Feed实例:
import time import feedparser from app import create_app, db from app.models import Feed # 替换成你的Feed模型类 def parallelParse(feed_url, feed_id): # 子进程内初始化App app = create_app() # 解析RSS/Atom d = feedparser.parse(feed_url) modified_parsed = d.feed.get('modified_parsed') if not modified_parsed: return # 没有修改时间就跳过 modified = time.mktime(modified_parsed) # 在App上下文中操作数据库 with app.app_context(): feed = Feed.query.get(feed_id) if feed and modified != feed.feedModified: feed.feedModified = modified db.session.commit()
第三步:在视图函数中调用多进程
只传递Feed的URL和ID,避免传递复杂对象:
from multiprocessing import Pool from flask import render_template from app.models import Feed from . import feeds @feeds.route('/') def home(): feeds_list = Feed.query.all() # 构造多进程任务参数 task_args = [(feed.url, feed.id) for feed in feeds_list] # 启动进程池执行任务 with Pool() as pool: pool.starmap(parallelParse, task_args) return render_template('home.html')
方案2:改用线程池(适合IO密集型场景)
如果你的RSS解析是IO密集型(大部分时间在等待网络请求),用线程比进程更轻量,而且线程共享内存空间,不需要pickle整个App对象。
修改视图函数和解析函数
from concurrent.futures import ThreadPoolExecutor from flask import current_app, render_template from app.models import Feed from . import feeds import time import feedparser def parallelParse(feed_url, feed_obj, app): d = feedparser.parse(feed_url) modified_parsed = d.feed.get('modified_parsed') if not modified_parsed: return modified = time.mktime(modified_parsed) if modified != feed_obj.feedModified: feed_obj.feedModified = modified with app.app_context(): db.session.commit() @feeds.route('/') def home(): feeds_list = Feed.query.all() # 获取真实的App对象(避免Flask的代理对象问题) app = current_app._get_current_object() with ThreadPoolExecutor(max_workers=4) as executor: # 提交任务 for feed in feeds_list: executor.submit(parallelParse, feed.url, feed, app) return render_template('home.html')
注意:Flask-SQLAlchemy的会话默认是线程本地的,所以线程中操作DB是安全的,但如果是自定义会话要额外注意线程安全。
为什么原来的代码会出错?
你之前直接把app传给子进程,而app内部关联了SQLAlchemy的shutdown_session局部函数——pickle序列化对象时会递归序列化所有关联的属性,碰到这个局部函数就失败了。而且跨进程共享Flask App本身就不符合设计规范,容易引发各种隐性问题。
内容的提问来源于stack exchange,提问作者Imran Said

