【问题标题】:Multithreading for Python DjangoPython Django 的多线程
【发布时间】:2013-08-24 16:32:54
【问题描述】:

一些功能应该在网络服务器上异步运行。发送电子邮件或数据后处理是典型的用例。

编写装饰器函数以异步运行函数的最佳(或最 Pythonic)方法是什么?

我的设置很常见:Python、Django、Gunicorn 或 Waitress、AWS EC2 标准 Linux

例如,这是一个开始:

from threading import Thread

def postpone(function):
    def decorator(*args, **kwargs):
        t = Thread(target = function, args=args, kwargs=kwargs)
        t.daemon = True
        t.start()
    return decorator

期望的用法:

@postpone
def foo():
    pass #do stuff

【问题讨论】:

  • 看看这个帖子stackoverflow.com/questions/573618/…。对于计划的作业,请选择基于 cron 的解决方案。 Scheduled Job, Asynchronous tasks 选择 Celery。在最近迁移到 Celery 之前,我从 github.com/tivix/django-cron 开始。
  • 感谢到目前为止的所有答案,但是 Celery 需要相当多的开销(安装应用程序,为其创建数据库)。因此,虽然 Celery 是一个解决方案,但它并没有回答我关于编写独立装饰器来多线程函数的问题。

标签: python django multithreading decorator python-multithreading


【解决方案1】:

我继续在规模和生产中使用此实现,没有任何问题。

装饰器定义:

def start_new_thread(function):
    def decorator(*args, **kwargs):
        t = Thread(target = function, args=args, kwargs=kwargs)
        t.daemon = True
        t.start()
    return decorator

示例用法:

@start_new_thread
def foo():
  #do stuff

随着时间的推移,堆栈不断更新和转换。

最初是 Python 2.4.7、Django 1.4、Gunicorn 0.17.2,现在是 Python 3.6、Django 2.1、Waitress 1.1。

如果你正在使用任何数据库事务,Django 将创建一个新的连接,这需要手动关闭:

from django.db import connection

@postpone
def foo():
  #do stuff
  connection.close()

【讨论】:

  • 我也在生产中运行相同的实现,没有任何问题,最好的部分是它可以与 uwsgi 一起使用,没有任何重大的性能问题
  • 请注意,这会泄漏数据库连接,因为 django 将为每个线程创建一个新的数据库连接,而您负责关闭它。
  • 测试确认@Chronial 是正确的。如果您的函数执行数据库事务(读取是我测试的),则创建一个新连接。然后线程终止连接仍然存在,并且在 80 后 Postgres 拒绝创建任何进一步的连接
  • 当服务器关闭并且你的延迟函数没有运行,或者正在运行的一部分时,它将被简单地中止。您需要从主线程调用新线程的join 方法,并且没有好的方法可以做到这一点。这也是人们使用 Celery 等的原因之一。
  • @Neil - 这完全取决于工作及其所做的工作 - 是否重要工作,是否使用数据库事务,是否会使某些内容处于不一致状态,是否有机制来检测并重试没有问题等。
【解决方案2】:

如果没有太多的新工作,tomcounsell 的方法很有效。如果在短时间内运行许多持久的作业,因此产生大量线程,主进程将受到影响。在这种情况下,您可以使用带有协程的线程池,

# in my_utils.py

from concurrent.futures import ThreadPoolExecutor

MAX_THREADS = 10


def run_thread_pool():
    """
    Note that this is not a normal function, but a coroutine.
    All jobs are enqueued first before executed and there can be
    no more than 10 threads that run at any time point.
    """
    with ThreadPoolExecutor(max_workers=MAX_THREADS) as executor:
        while True:
            func, args, kwargs = yield
            executor.submit(func, *args, **kwargs)


pool_wrapper = run_thread_pool()

# Advance the coroutine to the first yield (priming)
next(pool_wrapper)
from my_utils import pool_wrapper

def job(*args, **kwargs):
    # do something

def handle(request):
    # make args and kwargs
    pool_wrapper.send((job, args, kwargs))
    # return a response

【讨论】:

    【解决方案3】:

    Celery 是一个异步任务队列/作业队列。它有据可查,非常适合您的需要。我建议你开始here

    【讨论】:

      【解决方案4】:

      在 Django 中进行异步处理最常用的方法是使用Celerydjango-celery

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-03-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多