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

使用多进程时DataFrame已填充但获取为空的问题排查与解决

问题原因
  • 多进程内存隔离机制:Python多进程会复制父进程的内存空间,你传递给_populate_df子进程的self.df是原DataFrame的副本。子进程对副本的修改完全独立于主进程的原对象,因此外部访问ens.df时始终为空。
  • 未使用共享对象:你初始化了Manager()但未将self.df声明为进程共享对象,普通DataFrame不具备跨进程同步修改的能力。
修复方案

以下两种方案都能解决问题,按需选择:

方案1:用Manager创建共享容器收集数据(推荐)

通过Manager生成共享列表收集数据,最后在主进程转换为DataFrame,避免直接跨进程修改复杂对象的性能问题:

import numpy as np
import pandas as pd
from time import sleep
from multiprocessing import Process, Manager, Empty

class Ensambler():
    def __init__(self):
        self.manager = Manager()
        self.queue = self.manager.Queue(maxsize=20)
        self.df_queue = self.manager.Queue(maxsize=20)
        # 用共享列表存储数据,而非直接共享DataFrame
        self.data_list = self.manager.list()
        self.lock = self.manager.Lock()
            
    def _populate_queue(self):
        while True:
            try:
                c = np.random.randint(0,10)
                self.queue.put(c)
            except KeyboardInterrupt:
                break

    def _populate_df(self):
        while True:
            try:
                # 用Empty异常判断队列是否为空,比empty()更可靠
                item = self.df_queue.get(block=False)
                self.data_list.append(item)
                print(f"当前收集数据量:{len(self.data_list)}")
            except Empty:
                pass
            except KeyboardInterrupt:
                break
                
    def _process_queue(self,idx):
        while True:
            try:
                self.lock.acquire()
                item = self.queue.get()
                self.df_queue.put(item)
                self.lock.release()
                sleep(0.1)
            except KeyboardInterrupt:
                break
        
    
    def run(self,n_workers):
        try:
            populate = Process(target=self._populate_queue)
            populate_df = Process(target=self._populate_df)
            jobs = [populate,populate_df]
            for _ in range(n_workers):
                p = Process(target = self._process_queue, args=(_,))
                jobs.append(p)
            for job in jobs:
                job.start()

            # 等待用户中断(Ctrl+C)
            try:
                while True:
                    sleep(1)
            except KeyboardInterrupt:
                pass

            for job in jobs:
                job.join()
                job.close() if not job.is_alive() else job.terminate()

            # 主进程将共享列表转为DataFrame
            self.df = pd.DataFrame({'queue_items': self.data_list})

        except KeyboardInterrupt:
            for job in jobs:
                job.join()
                job.close() if not job.is_alive() else job.terminate()

if __name__ == '__main__':
    ens = Ensambler()
    ens.run(16)
    print(ens.df)

方案2:主进程直接处理队列生成DataFrame

取消_populate_df子进程,所有数据收集工作放在主进程完成,彻底规避跨进程对象修改问题:

import numpy as np
import pandas as pd
from time import sleep
from multiprocessing import Process, Manager, Empty

class Ensambler():
    def __init__(self):
        self.manager = Manager()
        self.queue = self.manager.Queue(maxsize=20)
        self.df_queue = self.manager.Queue(maxsize=20)
        self.df = pd.DataFrame()
        self.lock = self.manager.Lock()
            
    def _populate_queue(self):
        while True:
            try:
                c = np.random.randint(0,10)
                self.queue.put(c)
            except KeyboardInterrupt:
                break

    def _process_queue(self,idx):
        while True:
            try:
                self.lock.acquire()
                item = self.queue.get()
                self.df_queue.put(item)
                self.lock.release()
                sleep(0.1)
            except KeyboardInterrupt:
                break
        
    
    def run(self,n_workers):
        try:
            populate = Process(target=self._populate_queue)
            jobs = [populate]
            for _ in range(n_workers):
                p = Process(target = self._process_queue, args=(_,))
                jobs.append(p)
            for job in jobs:
                job.start()

            # 等待用户中断(Ctrl+C)
            try:
                while True:
                    sleep(1)
            except KeyboardInterrupt:
                pass

            for job in jobs:
                job.join()
                job.close() if not job.is_alive() else job.terminate()

            # 主进程从队列提取所有数据生成DataFrame
            items = []
            while True:
                try:
                    items.append(self.df_queue.get(block=False))
                except Empty:
                    break
            self.df = pd.DataFrame({'queue_items': items})

        except KeyboardInterrupt:
            for job in jobs:
                job.join()
                job.close() if not job.is_alive() else job.terminate()

if __name__ == '__main__':
    ens = Ensambler()
    ens.run(16)
    print(ens.df)
额外提示
  • 多进程中尽量避免直接共享复杂对象(如DataFrame),优先用队列、共享列表等轻量组件传递数据,最后在主进程组装,能减少跨进程同步的开销和潜在问题。
  • 不要依赖queue.empty()判断队列状态,该方法非进程安全,用try...except Empty捕获空队列更可靠。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 15:11:15