Celery和SQLAlchemy - 此结果对象不返回行.它已自动关闭

fed*_*qui 10 python multithreading sqlalchemy celery

我有一个连接到MySQL数据库的芹菜项目.其中一个表定义如下:

class MyQueues(Base):
    __tablename__ = 'accepted_queues'

    id = sa.Column(sa.Integer, primary_key=True)
    customer = sa.Column(sa.String(length=50), nullable=False)
    accepted = sa.Column(sa.Boolean, default=True, nullable=False)
    denied = sa.Column(sa.Boolean, default=True, nullable=False)
Run Code Online (Sandbox Code Playgroud)

此外,在我有的设置

THREADS = 4
Run Code Online (Sandbox Code Playgroud)

我被困在一个函数中code.py:

def load_accepted_queues(session, mode=None):

    #make query  
    pool = session.query(MyQueues.customer, MyQueues.accepted, MyQueues.denied)

    #filter conditions    
    if (mode == 'XXX'):
        pool = pool.filter_by(accepted=1)
    elif (mode == 'YYY'):
        pool = pool.filter_by(denied=1)
    elif (mode is None):
        pool = pool.filter(\
            sa.or_(MyQueues.accepted == 1, MyQueues.denied == 1)
            )

   #generate a dictionary with data
   for i in pool: #<---------- line 90 in the error
        l.update({i.customer: {'customer': i.customer, 'accepted': i.accepted, 'denied': i.denied}})
Run Code Online (Sandbox Code Playgroud)

运行时我收到一个错误:

[20130626 115343] Traceback (most recent call last):
  File "/home/me/code/processing/helpers.py", line 129, in wrapper
    ret_value = func(session, *args, **kwargs)
  File "/home/me/code/processing/test.py", line 90, in load_accepted_queues
    for i in pool: #generate a dictionary with data
  File "/home/me/envs/me/local/lib/python2.7/site-packages/sqlalchemy/orm/query.py", line 2341, in instances
    fetch = cursor.fetchall()
  File "/home/me/envs/me/local/lib/python2.7/site-packages/sqlalchemy/engine/base.py", line 3205, in fetchall
    l = self.process_rows(self._fetchall_impl())
  File "/home/me/envs/me/local/lib/python2.7/site-packages/sqlalchemy/engine/base.py", line 3174, in _fetchall_impl
    self._non_result()
  File "/home/me/envs/me/local/lib/python2.7/site-packages/sqlalchemy/engine/base.py", line 3179, in _non_result
    "This result object does not return rows. "
ResourceClosedError: This result object does not return rows. It has been closed automatically
Run Code Online (Sandbox Code Playgroud)

所以主要是它的一部分

ResourceClosedError: This result object does not return rows. It has been closed automatically
Run Code Online (Sandbox Code Playgroud)

有时也会出现这个错误:

DBAPIError :(错误)(,AssertionError('结果长度未请求长度:\n预期= 1.实际= 0.位置:21.数据长度:21',))'SELECT accepted_queues.customer AS accepted_queues_customer,accepted_queues.accepted AS accepted_queues_accepted ,accepted_queues.denied AS accepted_queues_denied \nFROM accepted_queues \nWHERE accepted_queues.accepted =%s OR accepted_queues.denied =%s'(1,1)

我无法正确地重现错误,因为它在处理大量数据时通常会发生.我试图改变THREADS = 4到1和错误消失了.无论如何,它不是一个解决方案,因为我需要保持线程的数量4.

另外,我对使用的必要性感到困惑

for i in pool: #<---------- line 90 in the error
Run Code Online (Sandbox Code Playgroud)

要么

for i in pool.all(): #<---------- line 90 in the error
Run Code Online (Sandbox Code Playgroud)

并且找不到合适的解释.

总之:任何建议跳过这些困难?

zzz*_*eek 12

总之:任何建议跳过这些困难?

是.你绝对不能同时在多个线程中使用Session(或与该Session相关的任何对象)或Connection,特别是对于DBAPI连接非常不安全的MySQL-Python*.您必须组织您的应用程序,以便每个线程处理它自己的专用MySQL-Python连接(以及因此与该Session相关联的SQLAlchemy Connection/Session /对象),而不会泄漏到任何其他线程.

  • 编辑:或者,您可以使用互斥锁将访问Session/Connection/DBAPI连接的权限限制为一次只有其中一个线程,尽管这种情况不太常见,因为所需的高度锁定往往会破坏使用目的首先是多个线程.

  • 如果你要使用线程,那么你需要掌握两种基本技术中的一种或两种 - 让每个线程根本不共享任何状态,或者使用共享的非线程安全资源的互斥锁/锁定一次只有一个线程触及它.使用SQLAclhemy,"不共享任何东西"用例通常使用[scoped_session]来解决(http://docs.sqlalchemy.org/en/rel_0_8/orm/session.html?highlight=scoped_session#unitofwork-contextual)构造.但是如果你牢牢掌握"线程局部变量"是什么以及它意味着什么,那么它就非常有用. (2认同)