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

如何用Python+asyncio测试AWS AppSync的GraphQL订阅?

解决AWS AppSync GraphQL订阅的异步测试问题

针对你在pytest-asyncio环境下测试AppSync订阅遇到的阻塞、就绪时机、任务管理问题,下面是适配测试流程的具体实现方案:

核心思路

  1. 用asyncio.create_task启动订阅任务,避免阻塞主线程
  2. 借助asyncio.Event标记订阅就绪状态,确保突变在订阅完全建立后执行
  3. 用列表收集订阅收到的消息,后续做断言验证
  4. 测试完成后通过任务取消机制终止订阅,避免资源泄漏

完整代码示例

import asyncio
import pytest
from gql import Client, gql
from gql.transport.aiohttp import AIOHTTPTransport

# 初始化AppSync客户端的fixture
@pytest.fixture
async def appsync_client():
    # 替换成你的AppSync API地址和认证信息
    transport = AIOHTTPTransport(
        url="https://your-api-id.appsync-api.region.amazonaws.com/graphql",
        headers={"Authorization": "your-auth-token"}
    )
    async with Client(transport=transport, fetch_schema_from_transport=False) as client:
        yield client

# 封装订阅执行逻辑,包含就绪通知和消息收集
async def run_subscription(client, ready_event, message_buffer):
    # 替换成你的订阅Query
    sub_query = gql("""
        subscription OnItemUpdated($itemId: ID!) {
            onItemUpdated(id: $itemId) {
                id
                name
                updatedAt
            }
        }
    """)
    # 启动订阅迭代,进入循环前订阅已就绪
    async for result in client.subscribe(sub_query, variable_values={"itemId": "test-001"}):
        # 首次迭代时触发就绪事件,通知主线程可以执行突变
        if not ready_event.is_set():
            ready_event.set()
        # 把收到的消息存入缓冲区
        message_buffer.append(result["onItemUpdated"])

# 订阅测试用例
@pytest.mark.asyncio
async def test_subscription_triggered_by_mutation(appsync_client):
    # 用于标记订阅就绪的事件
    subscription_ready = asyncio.Event()
    # 存储订阅收到的消息
    received_messages = []
    
    # 1. 启动订阅任务(非阻塞)
    sub_task = asyncio.create_task(
        run_subscription(appsync_client, subscription_ready, received_messages)
    )
    
    try:
        # 2. 等待订阅就绪,超时5秒避免无限等待
        await asyncio.wait_for(subscription_ready.wait(), timeout=5)
        
        # 3. 执行触发订阅的突变操作
        mutation_query = gql("""
            mutation UpdateItem($itemId: ID!, $newName: String!) {
                updateItem(id: $itemId, input: {name: $newName}) {
                    id
                    name
                }
            }
        """)
        await appsync_client.execute(
            mutation_query,
            variable_values={"itemId": "test-001", "newName": "Updated Test Item"}
        )
        
        # 4. 等待订阅收到消息,超时5秒
        await asyncio.wait_for(
            lambda: len(received_messages) > 0,
            timeout=5
        )
        
        # 5. 断言消息符合预期
        assert len(received_messages) == 1
        assert received_messages[0]["id"] == "test-001"
        assert received_messages[0]["name"] == "Updated Test Item"
        
    finally:
        # 6. 无论测试结果如何,取消订阅任务
        sub_task.cancel()
        # 捕获任务取消异常,避免报错
        try:
            await sub_task
        except asyncio.CancelledError:
            pass

关键细节说明

  • 订阅就绪通知:在订阅迭代器的第一次循环前,AppSync的连接已经建立完成,此时触发ready_event,确保突变不会提前执行
  • 任务管理:用asyncio.create_task创建订阅任务,主线程可以继续执行后续逻辑;测试结束后在finally块中取消任务,确保资源被正确释放
  • 超时控制:所有等待操作都加上asyncio.wait_for,防止因网络延迟或服务异常导致测试无限挂起
  • 消息收集:用列表作为消息缓冲区,异步订阅收到消息后存入列表,主线程可以直接访问做断言

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 15:22:48