Hola, estoy usando AsyncIOMotorClient para llamadas asíncronas de db a mongoDb. A continuación se muestra mi código.
xyz.py async def insertMany(self,collection_name,documents_to_insert): try: collection=self.database[collection_name] document_inserted = await collection.insert_many(documents_to_insert) return document_inserted except Exception: raise def insertManyFn(self,collection_name,documents_to_insert): try: loop=asyncio.new_event_loop() asyncio.set_event_loop(loop) loop1=asyncio.get_event_loop() inserted_documents_count = loop1.run_until_complete(self.insertMany(collection_name, documents_to_insert)) if inserted_documents_count==len(documents_to_insert): document_to_insert={Config.DB_JOB_COLUMN:Job.job_id,Config.DB_JOB_RESULT_COLUMN:Config.DB_JOB_RESULT_SUCCESS} loop1.run_until_complete(self.insertOne(Config.DB_JOB_COLLECTION, document_to_insert)) except Exception: raise xyz1.py t=Timer(10,xyz.insertManyFn,\ (collection_name,documents_to_insert)) t.start()Mientras ejecuto esto, obtengo una excepción.
RuntimeError: Task <Task pending coro=<xyz.insertMany() running at <my workspace location>/xyz.py:144> cb=[_run_until_complete_cb() at /usr/lib64/python3.5/asyncio/base_events.py:164]> got Future <Future pending cb=[_chain_future.<locals>._call_check_cancel() at /usr/lib64/python3.5/asyncio/futures.py:431]> attached to a different loopEn el programa anterior, se llamará a insertManyFn después de 10 segundos y se realizará la operación de inserción. Pero cuando hace la primera llamada a insertMany, obtengo una excepción.
Todavía quiero que mi MotorClient esté en el nivel superior del módulo, así que esto es lo que hago: parcheo MotorClient.get_io_loop para que siempre devuelva el bucle actual.
import asyncio import motor.core from motor.motor_asyncio import ( AsyncIOMotorClient as MotorClient, ) # MongoDB client client = MotorClient('mongodb://localhost:27017/test') client.get_io_loop = asyncio.get_running_loop # The current database ("test") db = client.get_default_database() # async context async def main(): posts = db.posts await posts.insert_one({'title': 'great success!') # Run main() asyncio.run(main())De acuerdo con la documentación , a AsyncIOMotorClient se le debe pasar un ioloop si no usa el predeterminado. Intente crear el cliente después de crear su ciclo de eventos:
loop=asyncio.new_event_loop() asyncio.set_event_loop(loop) client = AsyncIOMotorClient(io_loop=loop)He modificado el código y está funcionando.
def insertManyFn(self,loop,collection_name,documents_to_insert): try: inserted_documents_count = loop.run_until_complete(self.insertMany(event_loop,collection_name, documents_to_insert)) if len(inserted_documents_count)==len(documents_to_insert): document_to_insert={Config.DB_JOB_COLUMN:Job.job_id,Config.DB_JOB_RESULT_COLUMN:Config.DB_JOB_RESULT_SUCCESS} loop1.run_until_complete(self.insertOne(Config.DB_JOB_COLLECTION, document_to_insert)) except Exception: raise loop=asyncio.get_event_loop() t=Timer(10,self.xyz.insertManyFn,(loop,collection_name,documents_to_insert)) t.start()Explicación: estoy usando el temporizador de subprocesos de Python, que crea un subproceso propio para ejecutar una función después de un tiempo determinado. Entonces, dentro de este hilo estaba obteniendo un bucle de eventos que no debería ser el enfoque correcto, primero debería obtener el bucle de eventos y crear un hilo de temporizador en él. Supongo que esta es la única razón.