Python脚本如何实现进程间通信交互

wen python案例 29

本文目录导读:

Python脚本如何实现进程间通信交互

  1. 使用 multiprocessing 模块(最常用)
  2. 使用 socket(网络通信)
  3. 使用文件(File-based IPC)
  4. 使用 Manager(高级同步)
  5. 选择建议
  6. 完整示例:实时数据采集系统

Python 实现进程间通信(IPC, Inter-Process Communication)有多种方式,我会详细介绍最常用的几种方法,并给出完整的示例代码。

使用 multiprocessing 模块(最常用)

1 Queue(队列)

from multiprocessing import Process, Queue
import time
def producer(queue):
    """生产者进程"""
    for i in range(5):
        msg = f"消息 {i}"
        queue.put(msg)
        print(f"生产者发送: {msg}")
        time.sleep(1)
    queue.put("STOP")  # 发送结束信号
def consumer(queue):
    """消费者进程"""
    while True:
        msg = queue.get()
        if msg == "STOP":
            break
        print(f"消费者收到: {msg}")
        time.sleep(0.5)
if __name__ == "__main__":
    # 创建队列
    queue = Queue()
    # 创建并启动进程
    p1 = Process(target=producer, args=(queue,))
    p2 = Process(target=consumer, args=(queue,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

2 Pipe(管道)

from multiprocessing import Process, Pipe
import time
def sender(conn):
    """发送端"""
    messages = ["Hello", "World", "Python", "IPC", "END"]
    for msg in messages:
        conn.send(msg)
        print(f"发送: {msg}")
        time.sleep(0.5)
    conn.close()
def receiver(conn):
    """接收端"""
    while True:
        try:
            msg = conn.recv()
            if msg == "END":
                break
            print(f"接收: {msg}")
        except EOFError:
            break
    conn.close()
if __name__ == "__main__":
    # 创建管道 (duplex=True 表示双向通信)
    parent_conn, child_conn = Pipe(duplex=True)
    p1 = Process(target=sender, args=(parent_conn,))
    p2 = Process(target=receiver, args=(child_conn,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

3 Shared Memory(共享内存)

from multiprocessing import Process, Value, Array
import time
def increment(shared_value, shared_array):
    """修改共享数据"""
    for _ in range(10):
        # 修改共享值
        shared_value.value += 1
        # 修改共享数组
        shared_array[0] += 1
        shared_array[1] += 2
        print(f"Value: {shared_value.value}, Array: {shared_array[:]}")
        time.sleep(0.5)
def read_values(shared_value, shared_array):
    """读取共享数据"""
    for _ in range(5):
        print(f"读取 - Value: {shared_value.value}, Array: {shared_array[:]}")
        time.sleep(1)
if __name__ == "__main__":
    # 创建共享内存
    shared_value = Value('i', 0)  # 'i' 表示整数类型
    shared_array = Array('i', [0, 0, 0])  # 整数数组
    p1 = Process(target=increment, args=(shared_value, shared_array))
    p2 = Process(target=read_values, args=(shared_value, shared_array))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

使用 socket(网络通信)

import socket
import multiprocessing
import time
def server_process():
    """服务器进程"""
    # 创建UNIX域套接字
    server = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    server.bind('/tmp/test_socket')
    server.listen(1)
    conn, addr = server.accept()
    while True:
        data = conn.recv(1024)
        if not data:
            break
        print(f"服务器收到: {data.decode()}")
        conn.send(f"ACK: {data.decode()}".encode())
    conn.close()
    server.close()
def client_process():
    """客户端进程"""
    time.sleep(0.1)  # 等待服务器启动
    client = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
    client.connect('/tmp/test_socket')
    messages = ["Hello", "World", "END"]
    for msg in messages:
        client.send(msg.encode())
        response = client.recv(1024)
        print(f"客户端收到: {response.decode()}")
        time.sleep(0.5)
    client.close()
if __name__ == "__main__":
    import os
    try:
        os.unlink('/tmp/test_socket')
    except OSError:
        pass
    p1 = multiprocessing.Process(target=server_process)
    p2 = multiprocessing.Process(target=client_process)
    p1.start()
    p2.start()
    p1.join()
    p2.join()

使用文件(File-based IPC)

import time
import multiprocessing
import os
def writer_process(filename):
    """写入进程"""
    for i in range(5):
        with open(filename, 'w') as f:
            f.write(f"Message {i}\n")
        print(f"写入: Message {i}")
        time.sleep(1)
def reader_process(filename):
    """读取进程"""
    for _ in range(10):
        if os.path.exists(filename):
            with open(filename, 'r') as f:
                content = f.read().strip()
                if content:
                    print(f"读取: {content}")
        time.sleep(0.5)
if __name__ == "__main__":
    filename = "ipc_temp.txt"
    p1 = multiprocessing.Process(target=writer_process, args=(filename,))
    p2 = multiprocessing.Process(target=reader_process, args=(filename,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()
    # 清理临时文件
    if os.path.exists(filename):
        os.remove(filename)

使用 Manager(高级同步)

from multiprocessing import Process, Manager
import time
def worker_process(shared_dict, shared_list, lock):
    """工作进程"""
    with lock:
        shared_dict['count'] = shared_dict.get('count', 0) + 1
        shared_list.append(time.time())
        print(f"Worker added: {shared_dict['count']}")
if __name__ == "__main__":
    with Manager() as manager:
        # 创建共享数据结构
        shared_dict = manager.dict()
        shared_list = manager.list()
        lock = manager.Lock()
        processes = []
        for i in range(5):
            p = Process(target=worker_process, args=(shared_dict, shared_list, lock))
            processes.append(p)
            p.start()
        for p in processes:
            p.join()
        print(f"Final dict: {shared_dict}")
        print(f"Final list: {shared_list}")

选择建议

方法 适用场景 优点 缺点
Queue 生产者-消费者模式 线程安全,易用 需要序列化
Pipe 点对点通信 速度快,双向 仅两个进程
Shared Memory 大量数据共享 高效,无需序列化 复杂数据类型限制
Socket 跨机器通信 通用性强 需要序列化
Manager 复杂数据结构 支持多种数据类型 性能相对较慢

完整示例:实时数据采集系统

from multiprocessing import Process, Queue, Value, Manager
import random
import time
import signal
class DataCollector:
    """数据采集系统示例"""
    def __init__(self):
        self.data_queue = Queue()
        self.running = Value('b', True)  # 布尔类型共享变量
        self.stats = Manager().dict({
            'total_collected': 0,
            'error_count': 0
        })
    def collector(self, sensor_id):
        """数据采集进程"""
        while self.running.value:
            try:
                # 模拟采集数据
                data = {
                    'sensor_id': sensor_id,
                    'value': random.gauss(100, 15),
                    'timestamp': time.time()
                }
                self.data_queue.put(data)
                self.stats['total_collected'] += 1
                print(f"采集器 {sensor_id}: {data['value']:.2f}")
                time.sleep(random.uniform(0.5, 1.5))
            except Exception as e:
                self.stats['error_count'] += 1
                print(f"采集器 {sensor_id} 错误: {e}")
    def processor(self):
        """数据处理进程"""
        processed_count = 0
        while self.running.value or not self.data_queue.empty():
            try:
                data = self.data_queue.get(timeout=1)
                # 处理数据
                data['processed_value'] = data['value'] * 1.1
                data['processed'] = True
                processed_count += 1
                print(f"处理: 传感器 {data['sensor_id']} -> {data['processed_value']:.2f}")
            except:
                pass
        print(f"总处理数据: {processed_count}")
    def monitor(self):
        """监控进程"""
        while self.running.value:
            print(f"状态: 采集 {self.stats['total_collected']}, 错误 {self.stats['error_count']}")
            time.sleep(3)
    def start(self, duration=20):
        """启动系统"""
        # 创建3个采集进程
        collectors = [
            Process(target=self.collector, args=(i,))
            for i in range(3)
        ]
        # 创建处理进程
        processor = Process(target=self.processor)
        # 创建监控进程
        monitor = Process(target=self.monitor)
        # 启动所有进程
        for p in collectors + [processor, monitor]:
            p.start()
        # 运行指定时间
        time.sleep(duration)
        # 停止系统
        self.running.value = False
        # 等待所有进程结束
        for p in collectors + [processor, monitor]:
            p.join()
        print("系统停止")
        print(f"最终统计: {dict(self.stats)}")
if __name__ == "__main__":
    system = DataCollector()
    system.start(duration=15)

这些示例涵盖了Python中主要的进程间通信方式,选择哪种方式取决于你的具体需求:

  • 简单消息传递:使用 QueuePipe
  • 大量数据共享:使用 Shared Memory
  • 复杂数据结构:使用 Manager
  • 跨机器通信:使用 Socket
  • 实时性要求高:使用 PipeShared Memory

抱歉,评论功能暂时关闭!