本文目录导读:

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中主要的进程间通信方式,选择哪种方式取决于你的具体需求:
- 简单消息传递:使用
Queue或Pipe - 大量数据共享:使用
Shared Memory - 复杂数据结构:使用
Manager - 跨机器通信:使用
Socket - 实时性要求高:使用
Pipe或Shared Memory