RuntimeError: Task ___ running at ___ at got Future <Future pending cb=[Protocol._on_waiter_completed()]> attached to a different loop

Viewed 24

Faced a problem. I need to load information about objects and their properties into a database. I do this using a class with async methods that both use that class's (self.__session) attribute. If you call one of the methods without the other, then everything works as it should, and if you call them sequentially, then after executing one, the second throws a RuntimeError. I don't even know what to do. It is clear that the whole problem lies in the self.__session attribute, but I do not know how to fix it, probably I too green.

from asyncio import create_task, gather, run, set_event_loop_policy, WindowsSelectorEventLoopPolicy
from decimal import Decimal, ROUND_UP

from typing import Dict, List

from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker

from persistence.servises.validators import Product

from settings import POSTGRES
from models import ProductTable, PropertyValueTable

DATABASE = POSTGRES
database_name = DATABASE.get('DATABASE_NAME')
user = DATABASE.get('USERNAME')
password = DATABASE.get('PASSWORD')
async_driver = DATABASE.get('ASYNC_DRIVER')
host = DATABASE.get('HOST')
port = DATABASE.get('PORT')


class AsyncHandler:

    def __init__(self):
        self.__url = f'postgresql+{async_driver}://{user}:{password}@{host}:{port}/{database_name}'
        self.__engine = create_async_engine(self.__url, echo=False)
        self.__session = sessionmaker(self.__engine, expire_on_commit=False, class_=AsyncSession)

    async def add_product_to_db(self, key: int, product: dict) -> Dict[int, int]:
        async with self.__session() as session:
            prod = ProductTable(**product)
            session.add(prod)
            await session.commit()
        return {key: prod.id}

    async def add_value_to_db(self, property_value: dict):
        async with self.__session() as session:
            value = PropertyValueTable(**property_value)
            session.add(value)
            await session.commit()
        return value.id


async def add_product() -> Dict[int, int]:
    task_list = []
    for key in range(5):
        product = Product(category_id=33184,
                          name='some name',
                          slug='some slug',
                          description='some text',
                          price=Decimal(1234.56).quantize(Decimal('.01'), rounding=ROUND_UP),
                          photo1='1QdXEE_fxn2qmVYiR-81_RR5kUN5Ujm6j',
                          photo2=None,
                          photo3=None,
                          photo4=None)
        task_list.append(create_task(handler.add_product_to_db(key, product.dict())))
    result_list = list(await gather(*task_list))
    id_dict = {}
    for result in result_list:
        id_dict.update(result)
    return id_dict


async def add_property_value(id_dict: Dict[int, int], value_list: List[str]) -> List[int]:
    task_list = []
    n = 0
    for item_id in id_dict.values():
        task_list.append(create_task(handler.add_value_to_db({'product_id': item_id,
                                                              'property_id': 1177,
                                                              'value': value_list[n]}
                                                             )))
        n += 1
    return list(await gather(*task_list))


if __name__ == '__main__':
    value_list = ['some value 0', 'some value 1', 'some value 2', 'some value 3', 'some value 4']
    set_event_loop_policy(WindowsSelectorEventLoopPolicy())
    handler = AsyncHandler()

    items_id_dict = run(add_product())
    print(items_id_dict)

    values_id_list = run(add_property_value(items_id_dict, value_list))
    print(values_id_list)

Traceback:

File "C:/Work/parser/parser/persistence/servises/alchemy/asdasdasd.py", line 84, in <module>
    values_id_list = run(add_property_value(items_id_dict, value_list))
  File "C:\Users\user\AppData\Local\Programs\Python\Python38\lib\asyncio\runners.py", line 44, in run
    return loop.run_until_complete(main)
  File "C:\Users\user\AppData\Local\Programs\Python\Python38\lib\asyncio\base_events.py", line 616, in run_until_complete
    return future.result()
  File "C:/Work/parser/parser/persistence/servises/alchemy/asdasdasd.py", line 74, in add_property_value
    return list(await gather(*task_list))
  File "C:/Work/parser/parser/persistence/servises/alchemy/asdasdasd.py", line 41, in add_value_to_db
    await session.commit()
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\ext\asyncio\session.py", line 580, in commit
    return await greenlet_spawn(self.sync_session.commit)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\_concurrency_py3k.py", line 126, in greenlet_spawn
    result = context.throw(*sys.exc_info())
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 1451, in commit
    self._transaction.commit(_to_root=self.future)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 829, in commit
    self._prepare_impl()
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 808, in _prepare_impl
    self.session.flush()
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 3383, in flush
    self._flush(objects)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 3523, in _flush
    transaction.rollback(_capture_exception=True)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\langhelpers.py", line 70, in __exit__
    compat.raise_(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\compat.py", line 208, in raise_
    raise exception
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\session.py", line 3483, in _flush
    flush_context.execute()
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\unitofwork.py", line 456, in execute
    rec.execute(self)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\unitofwork.py", line 630, in execute
    util.preloaded.orm_persistence.save_obj(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\persistence.py", line 245, in save_obj
    _emit_insert_statements(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\orm\persistence.py", line 1238, in _emit_insert_statements
    result = connection._execute_20(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\base.py", line 1631, in _execute_20
    return meth(self, args_10style, kwargs_10style, execution_options)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\sql\elements.py", line 332, in _execute_on_connection
    return connection._execute_clauseelement(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\base.py", line 1498, in _execute_clauseelement
    ret = self._execute_context(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\base.py", line 1862, in _execute_context
    self._handle_dbapi_exception(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\base.py", line 2047, in _handle_dbapi_exception
    util.raise_(exc_info[1], with_traceback=exc_info[2])
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\compat.py", line 208, in raise_
    raise exception
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\base.py", line 1819, in _execute_context
    self.dialect.do_execute(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\engine\default.py", line 732, in do_execute
    cursor.execute(statement, parameters)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\dialects\postgresql\asyncpg.py", line 479, in execute
    self._adapt_connection.await_(
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\_concurrency_py3k.py", line 68, in await_only
    return current.driver.switch(awaitable)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\util\_concurrency_py3k.py", line 121, in greenlet_spawn
    value = await result
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\dialects\postgresql\asyncpg.py", line 408, in _prepare_and_execute
    await adapt_connection._start_transaction()
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\dialects\postgresql\asyncpg.py", line 716, in _start_transaction
    self._handle_exception(error)
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\dialects\postgresql\asyncpg.py", line 684, in _handle_exception
    raise error
  File "C:\Work\parser\venv\lib\site-packages\sqlalchemy\dialects\postgresql\asyncpg.py", line 714, in _start_transaction
    await self._transaction.start()
  File "C:\Work\parser\venv\lib\site-packages\asyncpg\transaction.py", line 138, in start
    await self._connection.execute(query)
  File "C:\Work\parser\venv\lib\site-packages\asyncpg\connection.py", line 318, in execute
    return await self._protocol.query(query, timeout)
  File "asyncpg\protocol\protocol.pyx", line 338, in query
RuntimeError: Task <Task pending name='Task-14' coro=<AsyncHandler.add_value_to_db() running at C:/Work/parser/parser/persistence/servises/alchemy/asdasdasd.py:41> cb=[gather.<locals>._done_callback() at C:\Users\user\AppData\Local\Programs\Python\Python38\lib\asyncio\tasks.py:766]> got Future <Future pending cb=[Protocol._on_waiter_completed()]> attached to a different loop
1 Answers
Related