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

Celery报RuntimeError: can't start new thread问题排查与解决求助

问题描述

我使用Celery(以Redis作为broker)通过Telegram Bot API发送消息。VPS配备4核共享CPU,将收件人列表按5人一组划分,每隔1-3秒向这些组发送消息。测试时(约2000名收件人),某个节点任务因RuntimeError: can't start new thread报错失败,此时CPU负载未超过15%,ulimit -u返回值为7795。

相关代码

models.py

class Post(models.Model):
    text = models.TextField(max_length=4096, blank=True, null=True, default=None, verbose_name="Text")

@receiver(post_save, sender=Post)
def instance_created(sender, instance, created, **kwargs):
    if created:
        pre_send_tg.apply_async((instance.id,), countdown=5)

tasks.py

def chunks(lst, n):
    res = []
    for i in range(0, len(lst), n):
        res.append(lst[i:i + n])
    return res


@celery_app.task(ignore_result=True)
def pre_send_tg(post_id):
    try:
        Post = apps.get_model('news.post')
        TelegramUser = apps.get_model('tools.telegramuser')

        post = Post.objects.get(id=post_id)
        users = [x.tg_id for x in TelegramUser.objects.all()]
            
        _start = datetime.datetime.now() + datetime.timedelta(seconds=5)
        count = 0

        for i in chunks(users, 5):
            for tg_id in i:          
                send_message.apply_async((tg_id, post.text), eta=_start)
          
            if count % 100 == 0:
                _start += datetime.timedelta(seconds=random.randint(50, 60))
            else:
                _start += datetime.timedelta(seconds=random.randint(2, 3))
            count += 1
    except Exception as e:
        print(e)


@celery_app.task(ignore_result=True, time_limit=10,  autoretry_for=(Exception,),
                retry_backoff=1800, retry_kwargs={'max_retries': 2})
def send_message(tg_id, text):
    bot = telebot.TeleBot(token)
    try:
        bot.send_message(chat_id=tg_id, text=text)
    except Exception as e:
        if e.args == (
                'A request to the Telegram API was unsuccessful. Error code: 403. Description: Forbidden: bot was blocked by the user',):
            pass
        elif e.args[0].startswith(
                "A request to the Telegram API was unsuccessful. Error code: 429."):
            raise Exception
        else:
            pass

Celery启动命令

celery -A Bot worker --loglevel=INFO --concurrency=10 -n worker1@%h --purge
celery -A Bot worker --loglevel=INFO --concurrency=10 -n worker2@%h --purge
celery -A Bot worker --loglevel=INFO --concurrency=10 -n worker3@%h --purge
问题原因
  1. Celery总并发过高:3个worker每个设置concurrency=10,总共有30个工作进程。每个进程处理send_message任务时,telebot内部(或依赖的requests库)会创建新线程处理网络请求,短时间内大量任务集中执行会导致线程数快速累积,触发进程的线程上限。
  2. TeleBot实例重复创建:每个send_message任务都新建telebot.TeleBot对象,该对象初始化时会创建线程相关资源(比如连接池线程),重复创建会加剧线程资源消耗。
  3. 任务集中执行:2000个收件人对应2000个send_message任务,这些任务按组设置了相同的eta时间,会在同一时间点集中执行,瞬间启动大量线程,超过单个进程的线程限制。
解决方法

1. 调整Celery并发配置

4核共享CPU建议总并发数控制在8以内,比如启动2个worker,每个concurrency=4:

celery -A Bot worker --loglevel=INFO --concurrency=4 -n worker1@%h --purge
celery -A Bot worker --loglevel=INFO --concurrency=4 -n worker2@%h --purge

2. 全局复用TeleBot实例

在tasks.py全局初始化Bot实例,避免重复创建:

# 全局初始化Bot,只创建一次
bot = telebot.TeleBot(token)

@celery_app.task(ignore_result=True, time_limit=10,  autoretry_for=(Exception,),
                retry_backoff=1800, retry_kwargs={'max_retries': 2})
def send_message(tg_id, text):
    try:
        bot.send_message(chat_id=tg_id, text=text)
    except Exception as e:
        if e.args == (
                'A request to the Telegram API was unsuccessful. Error code: 403. Description: Forbidden: bot was blocked by the user',):
            pass
        elif e.args[0].startswith(
                "A request to the Telegram API was unsuccessful. Error code: 429."):
            raise Exception
        else:
            pass

3. 分散任务执行时间

优化pre_send_tg的任务提交逻辑,给每个组的eta增加随机偏移,避免任务完全扎堆:

@celery_app.task(ignore_result=True)
def pre_send_tg(post_id):
    try:
        Post = apps.get_model('news.post')
        TelegramUser = apps.get_model('tools.telegramuser')

        post = Post.objects.get(id=post_id)
        users = [x.tg_id for x in TelegramUser.objects.all()]
            
        _start = datetime.datetime.now() + datetime.timedelta(seconds=5)
        count = 0

        for i in chunks(users, 5):
            # 给每组的eta增加0-1秒的随机偏移,分散执行时间
            group_eta = _start + datetime.timedelta(seconds=random.uniform(0, 1))
            for tg_id in i:          
                send_message.apply_async((tg_id, post.text), eta=group_eta)
          
            if count % 100 == 0:
                _start += datetime.timedelta(seconds=random.randint(50, 60))
            else:
                _start += datetime.timedelta(seconds=random.randint(2, 3))
            count += 1
    except Exception as e:
        print(e)

4. 改用线程池作为Celery任务池

将Celery的任务池从默认的prefork改为threads,减少进程开销,同时控制线程数:

celery -A Bot worker --loglevel=INFO --concurrency=8 --pool=threads -n worker1@%h --purge

5. 检查并调整进程线程限制

可以通过cat /proc/[worker_pid]/limits查看单个worker进程的线程数限制,如果确实是系统限制导致,可调整/etc/security/limits.conf增加线程限制(但优先通过代码优化解决):

* soft nproc 16384
* hard nproc 16384

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 16:55:17