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
问题原因
- Celery总并发过高:3个worker每个设置
concurrency=10,总共有30个工作进程。每个进程处理send_message任务时,telebot内部(或依赖的requests库)会创建新线程处理网络请求,短时间内大量任务集中执行会导致线程数快速累积,触发进程的线程上限。 - TeleBot实例重复创建:每个
send_message任务都新建telebot.TeleBot对象,该对象初始化时会创建线程相关资源(比如连接池线程),重复创建会加剧线程资源消耗。 - 任务集中执行: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
相关产品推荐
相关产品推荐

