Python多进程问题:接收AWS消息后对应进程未执行
问题排查:Python多进程配合AWS MQTT触发无响应
我正在完成大学的Python IoT模拟项目,需通过AWS消息触发对应进程执行代码。此前尝试线程方案,但因线程内代码量较大导致某线程崩溃,转而使用multiprocessing模块。
为复现问题,我编写了简化测试项目:接收AWS发送的1-4数值,触发四个进程中的对应进程打印指定内容,全局变量存储在variables1.py脚本中(如默认值为-1的var.process)。但运行测试项目后无任何反应,发送1-4消息后,对应进程未打印预期内容。完整项目将包含4个进程及AWS订阅、BLE通知两个中断场景,希望排查进程未执行的原因。
测试代码
from multiprocessing import Process from awscrt import mqtt, io from awsiot import mqtt_connection_builder import json import requests import boto3 from typing import Tuple from datetime import date import variables1 as var def Process_1(): while var.stop_thread == False: if (var.process == 1): print("Process 1 running...") def Process_2(): while var.stop_thread == False: if (var.process == 2): print("Process 2 running...") def Process_3(): while var.stop_thread == False: if (var.process == 3): print("Process 3 running...") def Process_4(): while (var.stop_thread == False): if (var.process == 4): print("Process 4 running...") #Crear los procesos process1 = Process(target=Process_1) process2 = Process(target=Process_2) process3 = Process(target=Process_3) process4 = Process(target=Process_4) event_loop_group = io.EventLoopGroup(1) host_resolver = io.DefaultHostResolver(event_loop_group) client_bootstrap = io.ClientBootstrap(event_loop_group, host_resolver) def on_message_received(topic, payload, dup, qos, retain, **kwargs): print("Message received") message = str(payload).replace("b'", "").replace("'", "") if (message == '1'): var.process = 1 elif (message == '1'): var.process = 2 elif (message == '1'): var.process = 3 elif (message == '1'): var.process = 4 ###################################################################### ################ FUNCTIONS FOR AWS CONFIGURATION... ################## # Cannot post this part of the code since it has private credentials # ###################################################################### #Iniciar los hilos process1.start() process2.start() process3.start() process4.start() while True: try: pass except KeyboardInterrupt: print("\nDispositivo desconectado") var.stop_thread = True process1.join() process2.join() process3.join() process4.join() break
问题根源分析
- 多进程全局变量共享失效:Python多进程拥有独立内存空间,主进程修改
var.process、var.stop_thread时,子进程中的这些变量是初始化时的副本,不会同步更新。子进程永远读不到主进程修改后的状态,自然不会触发打印逻辑。 - MQTT消息判断逻辑错误:
on_message_received函数中所有分支都判断message == '1',无论收到2、3还是4,都会把var.process设为1,逻辑完全错误。
修复方案
1. 修正MQTT消息处理逻辑
把条件判断改为对应数值,确保消息能正确映射到进程编号:
def on_message_received(topic, payload, dup, qos, retain, **kwargs): print("Message received") message = str(payload).replace("b'", "").replace("'", "") if message == '1': var.process = 1 elif message == '2': var.process = 2 elif message == '3': var.process = 3 elif message == '4': var.process = 4
2. 使用多进程通信机制共享状态
多进程无法直接共享全局变量,需用multiprocessing提供的同步原语实现状态同步,以下是两种可行方案:
方案一:用Manager共享变量
Manager可创建跨进程共享的字典/对象,实现主进程与子进程的状态同步:
from multiprocessing import Process, Manager import time def Process_1(shared_vars): while not shared_vars['stop_thread']: if shared_vars['process'] == 1: print("Process 1 running...") shared_vars['process'] = -1 # 触发后重置状态,避免重复打印 time.sleep(0.1) # 减少CPU占用 # 同理修改Process_2/3/4,接收shared_vars参数 if __name__ == '__main__': with Manager() as manager: shared_vars = manager.dict() shared_vars['stop_thread'] = False shared_vars['process'] = -1 # 创建进程时传入共享变量 process1 = Process(target=Process_1, args=(shared_vars,)) process2 = Process(target=Process_2, args=(shared_vars,)) process3 = Process(target=Process_3, args=(shared_vars,)) process4 = Process(target=Process_4, args=(shared_vars,)) # ... 保留原有AWS初始化代码 ... def on_message_received(topic, payload, dup, qos, retain, **kwargs): print("Message received") message = str(payload).replace("b'", "").replace("'", "") if message in ['1','2','3','4']: shared_vars['process'] = int(message) # 启动进程 process1.start() process2.start() process3.start() process4.start() while True: try: pass except KeyboardInterrupt: print("\nDispositivo desconectado") shared_vars['stop_thread'] = True process1.join() process2.join() process3.join() process4.join() break
方案二:用Queue传递触发指令
如果只需触发进程单次执行,无需持续检测状态,可通过队列传递任务指令:
from multiprocessing import Process, Queue def worker(queue, process_num): while True: task = queue.get() if task == 'stop': break if task == process_num: print(f"Process {process_num} running...") if __name__ == '__main__': task_queue = Queue() # 创建进程时传入队列和进程编号 process1 = Process(target=worker, args=(task_queue, 1)) process2 = Process(target=worker, args=(task_queue, 2)) process3 = Process(target=worker, args=(task_queue, 3)) process4 = Process(target=worker, args=(task_queue, 4)) # ... 保留原有AWS初始化代码 ... def on_message_received(topic, payload, dup, qos, retain, **kwargs): print("Message received") message = str(payload).replace("b'", "").replace("'", "") if message in ['1','2','3','4']: task_queue.put(int(message)) # 启动进程 process1.start() process2.start() process3.start() process4.start() while True: try: pass except KeyboardInterrupt: print("\nDispositivo desconectado") # 给所有进程发送停止指令 for _ in range(4): task_queue.put('stop') process1.join() process2.join() process3.join() process4.join() break
内容的提问来源于stack exchange,提问作者aag
相关产品推荐
相关产品推荐

