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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:47:31