MagicStack / MagicStack/asyncpg

asyncpg.exceptions._base.InterfaceError: cannot perform operation: another operation is in progress when executing gino processes in parallel

未关闭
#576 8 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看

还没有人认领这个 Issue。

主要语言
Python
星标
8.1k
派生
468
PR 合并指标
30 天内没有已合并 PR

描述

  • asyncpg version: 0.20.1
  • PostgreSQL version: 11
  • Do you use a PostgreSQL SaaS? If so, which? Can you reproduce
    the issue with a local PostgreSQL install?
    :no
  • Python version: 3.7
  • Platform: Win
  • Do you use pgbouncer?: no
  • Did you install asyncpg with pip?: yes
  • If you built asyncpg locally, which version of Cython did you use?:
  • Can the issue be reproduced under both asyncio and
    uvloop?
    : dont know, havent used uvloop

Hello! Thank you for your attention in advance.

Im trying to run parallel tasks in batches with asyncio, gino and asyncpg.

I have the magic starting from this entry point:

 while batch:
            tasks = {asyncio.create_task(JobService.run_job(job, connection_per_job=True, endpoint_connection=endpoint_request.db_connection)): job for job in batch}
            result = await asyncio.gather(*tasks)
            # executed_jobs.append({'name': job_instance.name, 'uuid': job_instance.uuid, 'result': job_execution.response})

            start += batch_jobs_count
            batch = jobs[start:start + batch_jobs_count]

endpoint_connection=endpoint_request.db_connection is the endpoint connection started from the point of receiving the request in my application (I'm using an internal util library over asyncpg, quart and gino)

Below in JobService.run_job I have:

@classmethod
    async def run_job(cls, job: Job, connection_per_job=False, endpoint_connection=None) -> (Job, JobExecution):
        print(f'received {job.name}, time: {datetime.utcnow()}')
        print(f'job: {job.name}, connection: {endpoint_connection.raw_connection._con._stmt_exclusive_section._acquired}')

        if connection_per_job and bool(endpoint_connection.raw_connection._con._stmt_exclusive_section._acquired):
            async with db_adapter.get_db().acquire() as conn:
                print('in context')
                job_instance, job_execution = await cls._run_job_internal(job, conn=conn)
        else:
            job_instance, job_execution = await cls._run_job_internal(job)

        return job_instance, job_execution

I'm trying to optimize using db connections by first checking if the general endpoint connection is free and using it if so and if not acquiring a new one in a context manager. In order to be sure which connection is used, I'm passing a bind (connection) parameter to all of my methods related to db operations (gino usage) and they seem to fail on the second created connection in the context, i don't really understand also how a connection created explicitly can be used in another operation already (from what i understand from the error)

贡献指南

这个仓库没有索引到贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

调研方向

Start with the asyncio.create_task/asyncio.gather entry point and JobService.run_job shown in the report, then reproduce the parallel batch with the same connection-per-job setup. Trace which connection each GINO operation uses and document a reproducible cause and safe connection-usage behavior; no repository file or test is named in the report.

由索引模型根据 Issue 内容生成。

评估

技术栈
postgresql, python
领域
backend, databases
Issue 类型
缺陷
难度
4/5
预计耗时
3-5 天
活跃度
停滞
描述清晰度
需要澄清
新手友好度
30/100

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。