基于asyncio Streams实现双向独立收发的Python TCP服务器方案咨询
可行性分析
Streams API完全能满足你的需求,它通过asyncio.start_server()创建的TCP连接可直接拿到reader/writer异步读写对象,结合asyncio的任务调度、锁机制,就能实现双向请求处理与并发控制,比Protocol API更直观易扩展。
核心实现思路
1. 连接会话与状态管理
- 为每个PLC连接维护独立会话对象,包含
reader、writer、当前是否有活跃数据流的标记(in_progress)、以及用于同步的asyncio.Lock()。 - 用全局字典保存在线PLC连接,以PLC标识(如IP+端口或自定义设备ID)为键,会话对象为值,方便快速查找。
2. 处理PLC主动发起的请求
为每个连接启动一个持续监听的异步任务,专门处理PLC主动发送的数据流:
import asyncio active_connections = {} class PLCSession: def __init__(self, reader, writer, plc_id): self.reader = reader self.writer = writer self.plc_id = plc_id self.in_progress = False self.lock = asyncio.Lock() async def plc_listener_task(session): while True: try: # 读取PLC请求(需根据实际协议处理粘包/拆包) raw_data = await session.reader.read(1024) if not raw_data: # 连接断开,清理会话 del active_connections[session.plc_id] session.writer.close() await session.writer.wait_closed() break # 检查是否有活跃数据流,有则拒绝 async with session.lock: if session.in_progress: session.writer.write(b"REJECT_ACK") await session.writer.drain() continue session.in_progress = True # 业务处理:解析PLC数据、写入RMQ等 processed_resp = process_plc_request(raw_data) # 发送响应并等待PLC结束ACK session.writer.write(processed_resp) await session.writer.drain() ack = await session.reader.read(32) if ack == b"PLC_ACK": async with session.lock: session.in_progress = False except Exception as e: # 异常后重置状态 async with session.lock: session.in_progress = False # 可选:记录日志、断开连接 print(f"PLC {session.plc_id} listener error: {e}") async def handle_plc_connection(reader, writer): addr = writer.get_extra_info('peername') plc_id = f"{addr[0]}:{addr[1]}" session = PLCSession(reader, writer, plc_id) active_connections[plc_id] = session # 启动监听任务 asyncio.create_task(plc_listener_task(session))
3. 处理RMQ触发的主动请求
单独启动RMQ消费任务,监听队列消息并触发向PLC的主动请求:
import aio_pika import json async def send_to_rmq_result(plc_id, data): # 向RMQ结果队列发送响应(示例) connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/") async with connection: channel = await connection.channel() await channel.default_exchange.publish( aio_pika.Message(body=json.dumps({"plc_id": plc_id, "data": data.decode()}).encode()), routing_key="plc_result_queue" ) async def rmq_consumer_task(): connection = await aio_pika.connect_robust("amqp://guest:guest@localhost/") async with connection: channel = await connection.channel() queue = await channel.declare_queue("plc_trigger_queue") async with queue.iterator() as queue_iter: async for message in queue_iter: async with message.process(): msg_content = json.loads(message.body) plc_id = msg_content["plc_id"] request_data = msg_content["request"].encode() # 检查PLC是否在线 if plc_id not in active_connections: continue session = active_connections[plc_id] # 检查并发流 async with session.lock: if session.in_progress: continue session.in_progress = True try: # 向PLC发送请求 session.writer.write(request_data) await session.writer.drain() # 等待PLC响应 resp_data = await session.reader.read(1024) # 转发结果到RMQ await send_to_rmq_result(plc_id, resp_data) # 发送结束ACK并等待PLC确认 session.writer.write(b"SERVER_ACK") await session.writer.drain() plc_ack = await session.reader.read(32) if plc_ack == b"PLC_ACK": async with session.lock: session.in_progress = False except Exception as e: async with session.lock: session.in_progress = False print(f"RMQ trigger error for PLC {plc_id}: {e}")
4. 主流程启动
同时启动TCP服务器和RMQ消费任务:
async def main(): server = await asyncio.start_server(handle_plc_connection, "0.0.0.0", 8888) async with server: # 启动RMQ消费任务 asyncio.create_task(rmq_consumer_task()) await server.serve_forever() if __name__ == "__main__": asyncio.run(main())
关键注意事项
- 并发安全:必须用
asyncio.Lock()保护in_progress状态,避免PLC监听任务和RMQ消费任务同时修改引发竞态。 - 协议适配:需根据PLC实际通信协议实现报文的粘包/拆包逻辑(比如固定长度、头部长度字段或分隔符),确保数据读取完整。
- 异常处理:针对连接断开、超时、读写失败等场景,及时清理会话、重置状态,避免无效资源占用。
- 任务管理:Python 3.11+可使用
asyncio.TaskGroup替代create_task,更优雅地批量管理异步任务的生命周期。
内容的提问来源于stack exchange,提问作者Alberto
相关产品推荐
相关产品推荐

