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

Python:如何从进程内的线程中提取数据?

多进程内线程的数据提取问题

我需要从进程内部的线程中提取数据,当前尝试启动多进程,每个进程内用线程执行任务,但遇到以下问题:

  • 直接用列表存储结果时,不使用mp.Process()代码正常,使用后无法获取数据
  • 尝试基于队列的写法,也没得到预期结果

原代码1:列表存储失败版本

import concurrent.futures as cf
import multiprocessing as mp
import subprocess
from dataclasses import dataclass

# Ping Output:
'''
Pinging google.com [142.250.113.100] with 32 bytes of data:
Reply from 142.250.113.100: bytes=32 time=26ms TTL=104
Reply from 142.250.113.100: bytes=32 time=26ms TTL=104
Reply from 142.250.113.100: bytes=32 time=27ms TTL=104
Reply from 142.250.113.100: bytes=32 time=26ms TTL=104

Ping statistics for 142.250.113.100:
    Packets: Sent = 4, Received = 4, Lost = 0 (0% loss),
Approximate round trip times in milli-seconds:
    Minimum = 26ms, Maximum = 27ms, Average = 26ms
'''

@dataclass(kw_only=True)
class Ping():
    name: str
    ip: str
    response: bool = False

def ping(target):
    cmd = subprocess.run(f"ping -n 2 {target}", shell=True, capture_output=True)
    out = cmd.stdout.decode()
    for line in out.splitlines():
        if "Reply from " in line:
            line = line.split(" ")
            input = Ping(name=f'{target}',ip=line[2][:-1], response=True)
            output.append(input) # Trying to append data to the output list.
            break

def thread(targets):
    with cf.ThreadPoolExecutor(max_workers=32) as execute:
        results = [execute.submit(ping, target) for target in targets]

targets = ["google.com", "amazon.com"]
output = []

if __name__ == "__main__":
    p = mp.Process(target=thread, args=(targets,))
    p.start()
    p.join()
    print(output) # This will not output any data from the ping function.

原代码2:队列写法失败版本

import concurrent.futures as cf
import multiprocessing as mp
import queue

def count(num, qq):
    qq.put(num)

def thread(numbers, q):
    qq = queue.Queue()
    with cf.ThreadPoolExecutor(max_workers=32) as execute:
        results = [execute.submit(count, num, qq) for num in numbers]
    q.put(qq.get())
    qq.task_done()

numbers = ['one', 'two', 'three']

if __name__ == "__main__":
    q = mp.Queue()
    p = mp.Process(target=thread, args=(numbers,q))
    p.start()
    p.join()
    print(q.get()) # This should contain 'one', 'two', 'three'.

问题根源

  1. 多进程拥有独立的内存空间,主进程的output列表和子进程内的output是完全不同的对象,子进程的修改不会同步到主进程
  2. 队列写法中错误混用了queue.Queue(仅适用于线程间通信)和mp.Queue(适用于进程间通信),且只向进程队列中放入了一个元素,未处理所有线程结果

修复方案

方案1:用mp.Queue实现进程间数据传递

把进程队列传递给子进程,线程任务直接将结果放入该队列,主进程最后从队列取出所有数据:

import concurrent.futures as cf
import multiprocessing as mp
import subprocess
from dataclasses import dataclass

@dataclass(kw_only=True)
class Ping():
    name: str
    ip: str
    response: bool = False

def ping(target, q):
    cmd = subprocess.run(f"ping -n 2 {target}", shell=True, capture_output=True)
    out = cmd.stdout.decode()
    for line in out.splitlines():
        if "Reply from " in line:
            line = line.split(" ")
            ping_result = Ping(name=f'{target}',ip=line[2][:-1], response=True)
            q.put(ping_result)  # 将结果放入进程队列
            break
    # 处理无响应的情况
    else:
        q.put(Ping(name=target, ip="", response=False))

def thread(targets, q):
    with cf.ThreadPoolExecutor(max_workers=32) as execute:
        # 给每个ping任务传递进程队列
        [execute.submit(ping, target, q) for target in targets]

targets = ["google.com", "amazon.com"]

if __name__ == "__main__":
    q = mp.Queue()
    p = mp.Process(target=thread, args=(targets, q))
    p.start()
    p.join()
    
    # 从队列取出所有结果
    output = []
    while not q.empty():
        output.append(q.get())
    print(output)

方案2:直接用线程池(IO密集型场景更高效)

ping属于IO密集型任务,GIL在IO等待时会自动释放,无需多进程,直接用线程池即可:

import concurrent.futures as cf
import subprocess
from dataclasses import dataclass

@dataclass(kw_only=True)
class Ping():
    name: str
    ip: str
    response: bool = False

def ping(target):
    cmd = subprocess.run(f"ping -n 2 {target}", shell=True, capture_output=True)
    out = cmd.stdout.decode()
    for line in out.splitlines():
        if "Reply from " in line:
            line = line.split(" ")
            return Ping(name=f'{target}',ip=line[2][:-1], response=True)
    # 处理无响应情况
    return Ping(name=target, ip="", response=False)

targets = ["google.com", "amazon.com"]

if __name__ == "__main__":
    output = []
    with cf.ThreadPoolExecutor(max_workers=32) as execute:
        # 用map批量获取线程结果
        results = execute.map(ping, targets)
        output.extend(results)
    print(output)

修复队列测试代码

直接将进程队列传递给线程任务,确保所有结果都放入进程队列:

import concurrent.futures as cf
import multiprocessing as mp

def count(num, q):
    q.put(num)

def thread(numbers, q):
    with cf.ThreadPoolExecutor(max_workers=32) as execute:
        # 直接传递进程队列给线程任务
        [execute.submit(count, num, q) for num in numbers]

numbers = ['one', 'two', 'three']

if __name__ == "__main__":
    q = mp.Queue()
    p = mp.Process(target=thread, args=(numbers,q))
    p.start()
    p.join()
    
    # 取出队列中所有元素
    output = []
    while not q.empty():
        output.append(q.get())
    print(output)  # 输出: ['one', 'two', 'three']

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 18:27:22