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'.
问题根源
- 多进程拥有独立的内存空间,主进程的
output列表和子进程内的output是完全不同的对象,子进程的修改不会同步到主进程 - 队列写法中错误混用了
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
相关产品推荐
相关产品推荐

