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

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操作可能因缓存失效导致处理速度下降,进一步加剧队列堆积。

优化建议

  1. 调整消费者任务获取逻辑:移除task_queue.get的timeout=0.5参数,改用阻塞式获取,避免消费者在队列有任务前的无效等待,减少CPU空闲时间。
  2. 限制队列最大容量:初始化JoinableQueue时设置maxsize参数,当队列满时生产者会自动阻塞,强制生产速度与消费速度匹配,防止无限制堆积。
  3. 使用共享内存存储bin数据:若需要共享bin的更新状态,改用multiprocessing.Manager或Array存储bin数据,避免每个消费者持有独立拷贝,减少内存开销和序列化操作。
  4. 实时监控速度指标:在生产者和消费者中添加计时逻辑,统计每秒生成的任务数和处理的任务数,明确两者的速度差异,精准定位瓶颈。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:04:51