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())
问题分析与修正方案
核心问题点
- 同步阻塞调用破坏异步逻辑:使用
time.sleep(0.1)会阻塞整个asyncio事件循环,导致异步任务get_nonblocking根本无法被调度执行。 - 任务执行顺序错误:创建任务后立刻检查
reply,此时任务还未开始执行,reply必然为空;回调方式也无法保证结果在预期时机返回。 - 任务取消未处理异常:直接调用
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
相关产品推荐
相关产品推荐

