Celery Worker配置需求:支持多线程执行但仅处理单个任务
我需要配置一个Celery Worker,使其能够使用多线程来运行本身包含多线程逻辑的Python程序,但不希望它接收并处理第二个任务。
我们有一个包含两个线程的Python程序,通过RabbitMQ/Celery触发执行。每台服务器都配备多处理器,我在每台服务器部署一个Worker,该Worker需要能够利用所有处理器资源。设置--concurrency=1时,额外的线程无法正常启动;而--concurrency设置大于1时,Worker会接收并处理额外的任务。
测试代码
启动第二个线程的代码:
import threading import time import os class im_watching(): def watcher_thread(self, done=False): with open(self.filen, 'a', encoding='utf-8') as f: f.write("I am watching\n") print("I am watching " + str(time.time())) x = 0 while not done: with open(self.filen, 'a', encoding='utf-8') as f: x = x+1 f.write("I am still watching\n") if x == 7: done = True with open(self.filen, 'a', encoding='utf-8') as f: f.write("I am done\n") def print_time(self): self.filen = 'threadfile_celery.txt' with open(self.filen, 'a', encoding='utf-8') as f: f.write('0: ' + str(time.time())+ '\n') print(time.time()) x = threading.Thread(target=self.watcher_thread, daemon=True) x.start() print(time.time()) r = range(10) for i in r: with open(self.filen, 'a', encoding='utf-8') as f: f.write(str(i+1)+ ' ' + str(time.time())+ '\n')
Celery任务调用
@app.task def check_thread(): import thread_watch as tw tp = tw.im_watching() tp.print_time()
运行对比
- 直接运行输出(正常并行):
0: 1660249134.2473516
1 1660249134.2533522
I am watching
2 1660249134.2653418
I am still watching
3 1660249134.268349
I am still watching
4 1660249134.272345
I am still watching
5 1660249134.2753484
6 1660249134.278356
I am still watching
7 1660249134.2803547
I am still watching
8 1660249134.2833555
I am still watching
9 1660249134.2993486
10 1660249134.3043566
I am still watching
I am done
- Celery默认配置运行输出(无并行):
0: 1660249011.9372554
I am watching
I am still watching
I am still watching
I am still watching
I am still watching
I am still watching
I am still watching
I am still watching
I am done
1 1660249011.9562619
2 1660249011.9582546
3 1660249011.960263
4 1660249011.9622533
5 1660249011.964253
6 1660249011.9652538
7 1660249011.9672582
8 1660249011.9692612
9 1660249011.971259
10 1660249011.9732592
使用以下命令启动Celery Worker:
celery -A your_app_name worker --concurrency=1 --pool=solo
参数说明
--concurrency=1:限制Worker同时处理的任务数量为1,确保不会接收并处理第二个任务--pool=solo:采用单线程执行池,这种池不会对任务内部的线程调度做限制。不同于默认的prefork池(即使--concurrency=1也会用子进程运行任务),solo池直接在Worker主线程中执行任务,任务内部的子线程可以和主线程正常并行调度,从而利用服务器的多处理器资源。
验证效果
启动Worker后调用check_thread任务,输出会和直接运行时的并行效果一致,同时Worker不会接收新的任务,直到当前任务执行完成。
内容的提问来源于stack exchange,提问作者AnonUser100000

