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

如何将Python WebSocket服务器改造为异步生成器?

基于WebSocket的异步爬虫生成器实现方案

先解决报错问题

你遇到的TypeError: 'Serve' object is not callable,大概率是把Serve类的实例当成函数调用了(比如写了serve()但没给类实现__call__方法),或者类的初始化/调用逻辑混乱。先理顺结构,再结合异步生成器做功能实现。

可行实现思路(无中间存储)

用Python的异步生成器结合aiohttp实现WebSocket服务器,把爬取到的HTML直接通过生成器yield返回,全程不需要中间数据库/文件存储。核心是把WebSocket的消息回调和异步生成器的产出逻辑绑定,收到爬取结果就立刻输出。

示例代码

import asyncio
from aiohttp import web

class WebSocketCrawler:
    def __init__(self):
        self._result_queue = asyncio.Queue()
        self._heartbeat_task = None

    # 异步生成器:对外返回爬取到的HTML
    async def crawl(self):
        while True:
            html = await self._result_queue.get()
            if html is None:  # 终止信号
                break
            yield html

    # WebSocket连接处理逻辑
    async def handle_websocket(self, request):
        ws = web.WebSocketResponse()
        await ws.prepare(request)

        # 启动心跳保活任务
        self._heartbeat_task = asyncio.create_task(self._send_heartbeat(ws))

        async for msg in ws:
            if msg.type == web.WSMsgType.TEXT:
                if msg.data == 'stop':
                    await ws.close()
                    self._heartbeat_task.cancel()
                    # 发送终止信号给生成器
                    await self._result_queue.put(None)
                else:
                    # 假设msg.data是爬取到的页面HTML
                    await self._result_queue.put(msg.data)
            elif msg.type == web.WSMsgType.CLOSED:
                self._heartbeat_task.cancel()
                await self._result_queue.put(None)
                break

        return ws

    # 心跳保活逻辑
    async def _send_heartbeat(self, ws):
        while True:
            await asyncio.sleep(30)
            if not ws.closed:
                await ws.send_str('heartbeat')

    # 启动WebSocket服务器
    async def start_server(self, host='0.0.0.0', port=8765):
        app = web.Application()
        app.add_routes([web.get('/ws', self.handle_websocket)])
        runner = web.AppRunner(app)
        await runner.setup()
        site = web.TCPSite(runner, host, port)
        await site.start()
        print(f"WebSocket server running on ws://{host}:{port}")

# 使用示例
async def main():
    crawler = WebSocketCrawler()
    # 启动服务器(后台运行)
    server_task = asyncio.create_task(crawler.start_server())
    # 遍历异步生成器获取HTML
    async for html in crawler.crawl():
        print(f"收到页面HTML,长度:{len(html)}")
        # 这里可以直接处理HTML,比如解析、存储等
    # 等待服务器关闭
    await server_task

if __name__ == '__main__':
    asyncio.run(main())

关键逻辑说明

  • 异步生成器crawl():通过异步队列_result_queue接收WebSocket传来的HTML,用yield返回给调用方,实现流式输出。
  • WebSocket处理:handle_websocket负责和爬虫工作端通信,收到HTML就丢进队列,收到停止指令就发送终止信号给生成器。
  • 心跳保活:单独起异步任务定时发送心跳帧,维持连接。
  • 无中间存储:HTML从WebSocket直接进入队列,再通过生成器输出,全程在内存流转,不需要额外存储介质。

为什么能解决你的问题

  1. 修正了"not callable"错误:类的方法都是正常调用,没有把实例当函数用的情况。
  2. 实现了异步返回HTML的生成器:调用方通过async for就能逐个获取爬取结果,符合复用模板的需求。
  3. 保留了心跳、指令控制等核心逻辑:WebSocket服务器的核心功能完整保留,还能轻松扩展其他指令。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 01:25:25