异步定时执行Elasticsearch查询仅输出sleeping无结果如何解决
问题根因
- 你使用了阻塞的
time.sleep(5),会卡住整个异步事件循环,导致你创建的execute_query异步任务完全没有被调度执行的机会,自然不会有结果输出 main函数是同步函数,不符合loop.run_until_complete()需要传入协程对象的要求,写法本身存在错误
修正后代码
from elasticsearch import AsyncElasticsearch import asyncio import random async def main(): es = AsyncElasticsearch("https://myElasticsearch", verify_certs=False, ssl_show_warn=False, request_cache=False) while True: asyncio.create_task(execute_query(es)) print("sleeping...") # 替换为异步sleep,不会阻塞事件循环,已创建的查询任务可以正常调度 await asyncio.sleep(5) async def execute_query(es): res = await es.search(index="myIndex", body = getRandomQuery()) print("I'm done") print(res) def getRandomQuery(): # 你的自定义随机查询逻辑 ... # Python3.7+推荐直接用asyncio.run,写法更简洁 asyncio.run(main()) # 如果你使用Python3.6及更低版本,保留原loop写法即可: # loop = asyncio.get_event_loop() # loop.run_until_complete(main())
效果说明
修改后完全符合你的预期:每隔5秒启动一个新的查询任务,新任务无需等待之前的查询执行完成,任意查询拿到返回结果后会立刻打印输出。
内容的提问来源于stack exchange,提问作者Mit94
相关产品推荐
相关产品推荐

