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
相关产品推荐
相关产品推荐

