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

Asyncio暂停/恢复:触发条件后先执行指定GraphQL请求再继续原订阅任务

实现方案

核心逻辑说明

你原代码中的print_handle是同步回调函数,无法直接执行异步的query2请求。要实现触发条件后暂停订阅任务、等待mutation请求完成再恢复的需求,仅需要把回调改为异步逻辑,触发条件时等待query2请求执行完毕即可——回调未返回时订阅任务会自动暂停处理后续消息,完全符合你的预期。

修改后的完整代码

from python_graphql_client import GraphqlClient
import asyncio
import os
import requests

headers={'Authorization': "2bTxxxxxxxxxxxxxxxxxxx"}

# 回调改为异步函数
async def print_handle(data):
    print(data["data"]["liveMeasurement"]["timestamp"]+" "+str(data["data"]["liveMeasurement"]["power"]))
    tall = (data["data"]["liveMeasurement"]["power"])
    if tall > 100:
        print("OK")
        # 等待mutation请求执行完成后再继续,自然暂停当前订阅处理
        await client.execute(
            query=query2, 
            headers={'Authorization': "2xxxxxxxxxxxxxxxx"}
        )

client = GraphqlClient(endpoint="wss://api.tibber.com/v1-beta/gql/subscriptions")

query = """
subscription{
  liveMeasurement(homeId:"xxxxxxxxxxxxxxxa"){
    timestamp
    power
  }
}
"""

query2 = """
mutation{
  sendPushNotification(input: {
    title: "Varsel! Høy belastning",
    message: "Du bruker nå høyere effekt enn 5 kw, pass på forbruket",
    screenToOpen: CONSUMPTION
  }){
    successful
    pushedToNumberOfDevices
  }
}
"""

async def main():
    # 异步订阅原生支持传入异步回调函数
    await client.subscribe(query=query, headers={'Authorization': "2xxxxxxxxxxxxxxxx"}, handle=print_handle)
    
asyncio.run(main())

兼容说明

  • 如果你使用的python_graphql_client版本subscribe不支持异步回调,可在原同步回调中调用asyncio.get_event_loop().run_until_complete(client.execute(...)),同样会阻塞回调直到请求完成后继续执行订阅逻辑。
  • 若client.execute为同步方法,可改为await asyncio.to_thread(client.execute, query=query2, headers={'Authorization': "2xxxxxxxxxxxxxxxx"})避免阻塞全局事件循环。
  • mutation请求执行完成后会自动释放对应连接,不影响原有订阅的持续运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 08:24:09