Python中Multiprocessing与共享数组功能异常问题排查
多进程版康威生命游戏显示异常问题
为了深入理解multiprocessing库,我用Python实现了一个利用多CPU核心的康威生命游戏。Pygame窗口能正确显示初始状态(存活和死亡细胞均匀分布),但工作进程修改后的所有细胞都显示为存活状态(暂时不确定实际值是否为1)。
我定位问题出在output_array[cell] = ...这行代码,尝试替换为output_array_shared[cell]、改用multiprocessing.Array并配合锁机制都没解决,只是让程序变慢了。以下是完整代码和最小复现示例:
完整代码
import math, multiprocessing from multiprocessing.sharedctypes import RawArray from queue import Empty import numpy as np width = 1600 height = 900 task_size_width = 64 task_size_height = 64 def update_canvas(screen, data): colors = np.array([(255, 0, 0), (0, 0, 255)], dtype=np.uint8) color_map = colors[data.reshape((height, width))] color_surface = pygame.surfarray.make_surface(color_map.swapaxes(0, 1)) screen.blit(color_surface, (0, 0)) pygame.display.flip() def live_neighbors(a, x, y): neighbors = [(x-1, y+1), (x, y+1), (x+1, y+1), (x-1, y), (x+1, y), (x-1, y-1), (x, y-1), (x+1, y-1)] live = 0 for x, y in neighbors: if 0 <= x < width and 0 <= y < height: if a[x + y * width] == 1: live += 1 return live def worker(queue, exit_flag, resume_exec, input_array_shared, output_array_shared, tasks_done): input_array = np.frombuffer(input_array_shared, dtype=np.uint8, count=width*height) output_array = np.frombuffer(output_array_shared, dtype=np.uint8, count=width*height) while not exit_flag.is_set(): try: x, y = queue.get_nowait() for i in range(task_size_height): for j in range(task_size_width): cx, cy = x * task_size_width + j, y * task_size_height + i if 0 <= cx < width and 0 <= cy < height: cell = cx + cy * width live = live_neighbors(input_array, cx, cy) if input_array[cell] == 1: if live == 2 or live == 3: output_array[cell] = 1 else: output_array[cell] = 0 else: if live == 3: output_array[cell] = 1 else: output_array[cell] = 0 with tasks_done.get_lock(): tasks_done.value += 1 except Empty: resume_exec.wait() print("worker: exiting") exit() if __name__ == "__main__": import os, time, pygame pygame.init() exit_flag = multiprocessing.Event() print("setup: setting up shared data") tasks = multiprocessing.Queue() resume_exec = multiprocessing.Event() input_array_shared = RawArray("B", width * height) output_array_shared = RawArray("B", width * height) #input_array_shared = multiprocessing.Array("B", width * height) #output_array_shared = multiprocessing.Array("B", width * height) input_array = np.frombuffer(input_array_shared, dtype=np.uint8, count=width*height) output_array = np.frombuffer(output_array_shared, dtype=np.uint8, count=width*height) total_tasks = 0 tasks_done = multiprocessing.Value("I", 0) tasks_done_copy = 0 print("setup: done!\n") print("setup: creating and starting processes") process_count = 12#os.cpu_count() processes = [multiprocessing.Process(target=worker, args=(tasks, exit_flag, resume_exec, input_array_shared, output_array_shared, tasks_done)) for _ in range(process_count)] for process in processes: process.start() print("setup: done!\n") time.sleep(0.1) print("setup: randomizing input state") input_array = np.random.randint(0, 2, size=width*height, dtype=np.uint8) print("setup: done!\n") print("window: initializing pygame") screen = pygame.display.set_mode((width, height)) pygame.display.set_caption(f"Conway's Game of Life - CPUs: {process_count}") print("window: done!\n") print("window: first refresh") output_array[:] = input_array[:] update_canvas(screen, output_array) print("window: done!\n") print("setup: preparing task order") tasks_ = [] tasks_w = math.ceil(width / task_size_width) tasks_h = math.ceil(height / task_size_height) tasks_cx = tasks_w // 2 tasks_cy = tasks_h // 2 tasks_x = tasks_cx tasks_y = tasks_cy d = 0 while d <= max(tasks_w, tasks_h): d += 1 if d == 1 and 0 <= tasks_x < tasks_w and 0 <= tasks_y < tasks_h: tasks_.append((tasks_x, tasks_y)) for i in range(d): tasks_y -= 1 if 0 <= tasks_x < tasks_w and 0 <= tasks_y < tasks_h: tasks_.append((tasks_x, tasks_y)) for i in range(d): tasks_x -= 1 if 0 <= tasks_x < tasks_w and 0 <= tasks_y < tasks_h: tasks_.append((tasks_x, tasks_y)) d += 1 for i in range(d): tasks_y += 1 if 0 <= tasks_x < tasks_w and 0 <= tasks_y < tasks_h: tasks_.append((tasks_x, tasks_y)) for i in range(d): tasks_x += 1 if 0 <= tasks_x < tasks_w and 0 <= tasks_y < tasks_h: tasks_.append((tasks_x, tasks_y)) print("setup: done!\n") while not exit_flag.is_set(): # reset the output array to all 0 print("mainloop: resetting output array") #output_array[:] = 0 # refill the task queue and reset the task counter print("mainloop: resetting queue") total_tasks = 0 for task in tasks_: tasks.put(task) total_tasks += 1 tasks_done.value = 0 tasks_done_copy = 0 # resume execution on the processes print("mainloop: resuming processes") resume_exec.set() time.sleep(0.1) resume_exec.clear() # keep updating the pygame window until all tasks are done processing while tasks_done_copy < total_tasks and not exit_flag.is_set(): #print("window: refresh") with tasks_done.get_lock(): tasks_done_copy = tasks_done.value update_canvas(screen, output_array) for event in pygame.event.get(): if event.type == pygame.QUIT: resume_exec.set() exit_flag.set() time.sleep(0.1) # an iteration has been computed, prepare for the next print("mainloop: preparing for next iteration\n") input_array[:] = output_array[:] print("main: waiting for workers to exit") for process in processes: process.join() print("main: exiting") os.system("taskkill /f /im python.exe")
最小复现示例
import math, multiprocessing, time from multiprocessing.sharedctypes import RawArray import numpy as np width = 1280 height = 720 def update_canvas(screen, data): colors = np.array([(255, 0, 0), (0, 0, 255)], dtype=np.uint8) color_map = colors[data.reshape((height, width))] color_surface = pygame.surfarray.make_surface(color_map.swapaxes(0, 1)) screen.blit(color_surface, (0, 0)) pygame.display.flip() def worker(output_array_shared, _): output_array = np.frombuffer(output_array_shared, dtype=np.uint8, count=width*height) for y in range(height): for x in range(width): output_array[x + y * width] = 0 time.sleep(0.0001) if __name__ == "__main__": import pygame pygame.init() output_array_shared = RawArray("B", width * height) output_array = np.frombuffer(output_array_shared, dtype=np.uint8, count=width*height) screen = pygame.display.set_mode((width, height)) output_array[:] = np.random.randint(0, 2, size=width*height, dtype=np.uint8)[:] process = multiprocessing.Process(target=worker, args=(output_array_shared, 0)) process.start() while True: update_canvas(screen, output_array) time.sleep(0.1)
内容的提问来源于stack exchange,提问作者PFnove
相关产品推荐
相关产品推荐

