如何用Asyncio实现非阻塞多客户端TCP请求?
问题分析与解决方案
你遇到的核心问题是错误地使用了loop.run_until_complete(),导致每次调用main()都会阻塞线程,无法实现多客户端并发,尤其是当发送'1'时,read_position是无限循环协程,会让第一个main()永远不会结束,后续代码完全无法执行。
问题根源拆解
loop.run_until_complete(coroutine)的作用是阻塞当前线程,直到传入的协程完全执行完毕。你先调用main('2'),程序会等这个客户端的所有操作(发消息、收响应、关连接)做完,才会执行下一行的main('3')。- 当消息是
'1'时,read_position是无限循环的协程,run_until_complete会一直等待它结束,导致第一个main()永远挂起,后续客户端根本没有启动机会。
解决方案一:普通脚本实现多客户端并发
我们需要把多个客户端协程放到同一个事件循环中,用asyncio.gather来并发执行,避免每次调用都阻塞线程:
import asyncio PYTHONASYNCIODEBUG = 1 # ECHO CLIENT PROTOCOL async def tcp_echo_client(message, loop): reader, writer = await asyncio.open_connection('127.0.0.1', 8888, loop=loop) print('Send: %r' % message) writer.write(message.encode()) if message == '1': await read_position(reader) else: await read_server(reader) print('Close the socket') writer.close() await writer.wait_closed() # 确保连接完全关闭 async def read_server(reader): server_message = await reader.read() print(type(server_message)) print('Received: %r' % server_message.decode()) async def read_position(reader): try: while True: print("I'm Here") server_message = await reader.read(50) if not server_message: # 服务器关闭连接时退出循环,避免无限挂起 break position = server_message.split() print(position) print(type(position)) print('Received: %r' % server_message.decode()) finally: print("Position reader exiting") async def main(): # 创建多个客户端协程,用gather并发执行 tasks = [ tcp_echo_client('2', asyncio.get_event_loop()), tcp_echo_client('3', asyncio.get_event_loop()), # 可以添加更多客户端任务 ] await asyncio.gather(*tasks) if __name__ == '__main__': asyncio.run(main()) # 用asyncio.run简化事件循环管理
解决方案二:适配Kivy应用的并发(重点)
Kivy有自己的主UI事件循环,不能直接用asyncio.run或run_until_complete阻塞主线程(否则UI会卡死),需要把asyncio任务提交到后台运行:
from kivy.app import App from kivy.uix.button import Button from kivy.uix.boxlayout import BoxLayout from kivy.clock import Clock import asyncio # 复用上述的tcp_echo_client、read_server、read_position函数 async def tcp_echo_client(message, loop): reader, writer = await asyncio.open_connection('127.0.0.1', 8888, loop=loop) print('Send: %r' % message) writer.write(message.encode()) if message == '1': await read_position(reader) else: await read_server(reader) print('Close the socket') writer.close() await writer.wait_closed() async def read_server(reader): server_message = await reader.read() msg = server_message.decode() print('Received: %r' % msg) # 如果需要更新UI,必须用Clock调度到Kivy主线程 # Clock.schedule_once(lambda dt: self.update_display(msg)) async def read_position(reader): try: while True: print("I'm Here") server_message = await reader.read(50) if not server_message: break position = server_message.split() print(position) print('Received: %r' % server_message.decode()) finally: print("Position reader exiting") class RobotControlApp(App): def build(self): self.loop = asyncio.get_event_loop() layout = BoxLayout(orientation='vertical') btn_task1 = Button(text="Send '1' (Position Task)") btn_task1.bind(on_press=self.send_task1) btn_task2 = Button(text="Send '2' (General Task)") btn_task2.bind(on_press=self.send_task2) layout.add_widget(btn_task1) layout.add_widget(btn_task2) return layout def send_task1(self, instance): # 把协程提交到asyncio事件循环,不阻塞Kivy UI线程 asyncio.create_task(tcp_echo_client('1', self.loop)) def send_task2(self, instance): asyncio.create_task(tcp_echo_client('2', self.loop)) def on_stop(self): # 应用关闭时清理事件循环 self.loop.close() if __name__ == '__main__': RobotControlApp().run()
额外提示
- 如果需要在Kivy中更新UI(比如把服务器返回的位置信息显示到界面),不能直接在asyncio协程中操作UI组件,必须用
Clock.schedule_once把UI更新任务放到Kivy的主循环中执行。 - 确保你的服务器是异步多连接的(比如用
asyncio.start_server创建的服务器),这样才能同时处理多个客户端的请求。
内容的提问来源于stack exchange,提问作者Filipe Ribeiro
相关产品推荐
相关产品推荐

