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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 14:05:14