Python多进程生产者消费者队列持续增长问题排查
业务场景
生产者通过迭代哈希生成数值,将其哈希映射至多个bin,每个数值会被放入两个bin中;每个消费者负责一组bin的地址范围,生产者将数值放入对应消费者的工作队列,消费者从队列取出任务执行powmod(模幂运算)后更新bin数据。
问题现象
测试生产者与消费者的配比时发现:1:7的配比能维持队列大小稳定,但2:14的配比却无法阻止队列线性增长;即使增加消费者数量,仅能减缓队列增长速度,无法彻底停止增长。在云环境中可确保每个进程独占CPU核心。
测试与尝试
- 固定任务量测试显示:消费者数量加倍,处理耗时减半,符合预期;按此推算2:16的配比应能平衡生产与消费,但实际运行5分钟后,队列仍增长至45-47万条。
- 已验证任务分配均匀、队列任务可正常移除;尝试使用concurrent futures进程池与Manager管理队列,虽能同步完成任务,但CPU负载仅30%,任务处理耗时更长,而当前实现下CPU负载达100%。
核心代码
import hashlib import math import sys import time from multiprocessing import Process, JoinableQueue, cpu_count from gmpy2 import mpz, powmod, is_prime, bit_set from random import randint, getrandbits from secrets import token_urlsafe from queue import Empty class Blake2b(object): def __init__(self, output_bit_length=64, salt=None, optimized64=False): self.output_bit_length = output_bit_length digest_size = math.ceil(output_bit_length / 8) # convert to bytes if salt is None: self.hash = hashlib.blake2b(digest_size=digest_size) else: self.hash = hashlib.blake2b(digest_size=digest_size, salt=salt) self.output_modulus = mpz(2**(self.output_bit_length - 1)) # integer double the size of the bit length - 1 (so double -2) def generating_value(self, preimage): h = self.hash.copy() h.update(preimage) current_digest = h.digest() while True: U = int.from_bytes(current_digest, sys.byteorder) % self.output_modulus candidate = bit_set(self.output_modulus + U, 0) if is_prime(candidate): return candidate h.update(b'0') current_digest = h.digest() def KeyGen(self, preimage, size): h = self.hash.copy() h.update(preimage) digest = h.hexdigest() return int(digest, 16) % size def task(bins, task_queue, ident, f): while True: try: next_task = task_queue.get(timeout=0.5) except Empty: print('Consumer: gave up waiting...', file=f, flush=True) continue if next_task is None: # Poison pill tells it to shutdown print (f'Consumer {ident}: Exiting', flush=True, file=f) task_queue.task_done() break value, in_bin = next_task N, A = bins[in_bin] answer = powmod(A, value, N) bins[in_bin][1] = answer task_queue.task_done() return def add_to_queue(tableSize, num_consumers, task_queues, f): bit_length = 86 #The security parameter b = Blake2b(output_bit_length=bit_length) leftKey = Blake2b(output_bit_length=10, salt=b'left') rightKey = Blake2b(output_bit_length=10, salt=b'right') bins_per_process = tableSize // num_consumers while True: rnd_val = b.generating_value(token_urlsafe(16).encode('utf-8')) KeyLeft = leftKey.KeyGen(str(rnd_val).encode('utf-8'), tableSize) KeyRight = rightKey.KeyGen(str(rnd_val).encode('utf-8'), tableSize) task_queues[KeyLeft//bins_per_process].put([rnd_val, KeyLeft % bins_per_process]) task_queues[KeyRight//bins_per_process].put([rnd_val, KeyRight % bins_per_process]) return if __name__ == '__main__': num_cpus = cpu_count() num_consumers = 16 num_producers = 2 table_size = num_consumers*8 #how many bins total with open("qGrowthTests.txt", "a") as f: table = [[getrandbits(3072), randint(2, 2**3072 - 1)] for _ in range(table_size)] bins_lst = [table[(i * len(table)) // num_consumers:((i + 1) * len(table)) // num_consumers] for i in range(num_consumers)] task_queues = [JoinableQueue() for _ in range(num_consumers)] producers = [] for i in range(num_producers): producer = Process(target=add_to_queue, args=(table_size, num_consumers, task_queues, f)) producers.append(producer) producer.start() consumers = [] for queue_num, queue in enumerate(task_queues, start=0): process = Process(target=task, args=(bins_lst[queue_num], queue, queue_num, f)) consumers.append(process) process.start() # checking the size of the queues once every 5 seconds for a 5 minute test n = 1 while n < 60: # 5 minutes = 300 seconds, 300/5=60 time.sleep(5) for q_num, q in enumerate(task_queues, start=0): print("Queue ", q_num, " has: ", q.qsize(), file=f) n += 1 # Force all the processes to finish for process in producers: if process.is_alive(): print (f'Producer {process}: Exiting', flush=True, file=f) process.terminate() process.join() for process in consumers: if process.is_alive(): print (f'Consumer {process}: Exiting', flush=True, file=f) process.terminate() process.join() print('Test Complete')
问题原因分析
1. 生产者任务生成速度的非线性扩展
单个生产者运行时,队列put操作的进程间通信开销会占据部分CPU时间,导致实际生成速度未达单个核心的极限。当启动多个生产者时,每个生产者的通信开销占比下降,总生成速度会超过单个生产者的N倍(N为生产者数量),而消费者的处理能力仅为线性扩展,无法匹配突增的生产速度。
2. 消费者任务处理的隐性开销
固定任务量测试仅验证了纯处理能力的线性扩展,但实时运行时,消费者需要频繁从队列获取任务、反序列化任务数据,这些操作在队列堆积时会产生额外开销(如大整数序列化延迟),导致实际处理速度低于预期。此外,当前消费者修改的是bin的本地拷贝(非共享内存),无意义的内存操作可能增加隐性开销。
3. 任务处理时间的波动性
powmod的处理时间依赖于模数(3072位)和指数(86位)的具体值,存在显著波动。当出现一批耗时极长的任务时,队列会瞬间堆积,而生产者持续生成任务,即使平均处理能力足够,也会导致队列持续增长。
4. 队列实现的效率瓶颈
JoinableQueue基于管道和锁实现,当队列规模过大时,任务的序列化/反序列化开销会显著增加,同时消费者的get操作可能因缓存失效导致处理速度下降,进一步加剧队列堆积。
优化建议
- 调整消费者任务获取逻辑:移除
task_queue.get的timeout=0.5参数,改用阻塞式获取,避免消费者在队列有任务前的无效等待,减少CPU空闲时间。 - 限制队列最大容量:初始化
JoinableQueue时设置maxsize参数,当队列满时生产者会自动阻塞,强制生产速度与消费速度匹配,防止无限制堆积。 - 使用共享内存存储bin数据:若需要共享bin的更新状态,改用
multiprocessing.Manager或Array存储bin数据,避免每个消费者持有独立拷贝,减少内存开销和序列化操作。 - 实时监控速度指标:在生产者和消费者中添加计时逻辑,统计每秒生成的任务数和处理的任务数,明确两者的速度差异,精准定位瓶颈。
内容的提问来源于stack exchange,提问作者LLL

