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

Python asyncio非阻塞TCP客户端监听第三方服务遇阻求助

问题:非阻塞监听TCP服务数据异常

我需要监听一个无法控制的本地Windows TCP服务,该服务每500ms通过固定套接字广播数据。阻塞模式下能正常读取随机长度的数据包,但用asyncio封装get_one_line()实现非阻塞监听时出问题:移除10次循环限制后仅打印count无数据,保留限制则因任务残留报错。

阻塞模式可正常运行的代码

import asyncio, telnetlib3, time
TCP_SERVER_ADDRESS = "127.0.0.7"
TCP_SERVER_PORT = 8000

async def get_one_line(reader):
    reply1 = []
    while True:
        c = await reader.read(1)
        if not c:
            break
        if c in ['\r', '\n']:
            break
        reply1.append(c)
    return reply1

async def main():
    reader, writer = await telnetlib3.open_connection(TCP_SERVER_ADDRESS, TCP_SERVER_PORT)
    writer.write("$SYS,INFO")
    count = 0
    while True:
        reply = []
        reply = await get_one_line(reader)
        if reply:
            print('reply:', ''.join(reply))
            del reply[:]
        print(count)
        count += 1
        if count > 10:
            break
        time.sleep(0.1)

asyncio.run(main())

非阻塞尝试的问题代码

import asyncio, telnetlib3, time
TCP_SERVER_ADDRESS = "127.0.0.7"
TCP_SERVER_PORT = 8000

async def get_one_line(reader):
    reply1 = []
    while True:
        c = await reader.read(1)
        if not c:
            break

        if c in ['\r', '\n']:
            break

        reply1.append(c)
    return reply1

async def get_nonblocking(reader, callback):
    reply2 = await get_one_line(reader)
    callback(reply2)

async def main():
    reader, writer = await telnetlib3.open_connection(TCP_SERVER_ADDRESS, TCP_SERVER_PORT)
    writer.write("$SYS,INFO")
    count = 0
    while True:
        reply = []
        getTask = asyncio.create_task(get_nonblocking(reader,reply.append))
        if reply:
            print('reply:', ''.join(reply))
            del reply[:]
        print(count)
        count += 1
        if count > 10:
            getTask.cancel()
            break
        time.sleep(0.1)

asyncio.run(main())

问题分析与修正方案

核心问题点

  1. 同步阻塞调用破坏异步逻辑:使用time.sleep(0.1)会阻塞整个asyncio事件循环,导致异步任务get_nonblocking根本无法被调度执行。
  2. 任务执行顺序错误:创建任务后立刻检查reply,此时任务还未开始执行,reply必然为空;回调方式也无法保证结果在预期时机返回。
  3. 任务取消未处理异常:直接调用getTask.cancel()后未等待任务结束、捕获取消异常,会导致程序报错。

修正后的代码

import asyncio, telnetlib3
TCP_SERVER_ADDRESS = "127.0.0.7"
TCP_SERVER_PORT = 8000

async def get_one_line(reader):
    reply1 = []
    while True:
        c = await reader.read(1)
        if not c:
            break
        if c in ['\r', '\n']:
            break
        reply1.append(c)
    return reply1

async def listen_data(reader):
    # 独立的监听任务,持续读取并处理数据
    while True:
        reply = await get_one_line(reader)
        if reply:
            print('reply:', ''.join(reply))

async def main():
    reader, writer = await telnetlib3.open_connection(TCP_SERVER_ADDRESS, TCP_SERVER_PORT)
    writer.write("$SYS,INFO")
    
    # 启动异步监听任务
    listen_task = asyncio.create_task(listen_data(reader))
    
    count = 0
    try:
        while True:
            print(count)
            count += 1
            # 用asyncio.sleep替代time.sleep,不阻塞事件循环
            await asyncio.sleep(0.1)
            # 移除循环限制可删除此判断
            if count > 10:
                break
    finally:
        # 正确取消任务:先取消,再等待任务结束并捕获取消异常
        listen_task.cancel()
        try:
            await listen_task
        except asyncio.CancelledError:
            print("监听任务已取消")

asyncio.run(main())

修正说明

  • 将数据监听逻辑拆分为独立异步任务listen_data,与count打印逻辑并行执行,真正实现非阻塞。
  • 替换time.sleep为await asyncio.sleep,确保事件循环能正常调度所有异步任务。
  • 任务取消时,通过await listen_task并捕获CancelledError,避免残留任务导致报错。
  • 直接在监听任务内处理数据打印,无需回调传递结果,逻辑更简洁可靠。

内容的提问来源于stack exchange,提问作者ilnadi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 17:33:26