MagicStack / MagicStack/asyncpg

Listener callback that accesses database does not complete fetch

オープン
#815 コメント 3 件 リアクション 0 件 担当者 0 名 GitHub で見る

まだ誰も着手していません。

主要言語
Python
スター
8.1k
フォーク
468
PR マージ指標
30日以内にマージされた PR はありません

説明

* **asyncpg version**: 0.24.0
* **PostgreSQL version**: 12.8
* **Do you use a PostgreSQL SaaS? If so, which? Can you reproduce
the issue with a local PostgreSQL install?**: Local PostgreSQL install
* **Python version**: 3.9.1
* **Platform**: Linux Mint
* **Do you use pgbouncer?**: No
* **Did you install asyncpg with pip?**: Yes
* **If you built asyncpg locally, which version of Cython did you use?**: N/A
* **Can the issue be reproduced under both asyncio and
[uvloop](https://github.com/magicstack/uvloop)?**: Have not tested with uvloop yet

When a row in my table is inserted or updated, I need a listener to pull data from the table and perform operations with it. The callback is successfully triggered on updates and inserts, but **within the callback the database cannot be accessed.** A new connection is successfully acquired from the pool (line [1]) but **the fetch operation does not complete** (line [2]).

Is this an issue with my code, or is it a limitation of asyncpg?

Sorry for an enormous code example - it's in a microservice application with a number of Python applications that all access the database so this was about as minimal I could make things while keeping it vaguely resembling the existing program structure.

```python
import asyncio
import asyncpg
from os import getenv
from dotenv import load_dotenv
import json

load_dotenv("../../.env")

class DatabaseInterface:
"""
Multiple different applications need to access the database, and will be performing similar operations.
Methods and properties to do with interacting with the database are all put into the DatabaseInterface class,
with each application creating a new instance of the class.
"""
def __init__(self):
self.pool = None
self.listeners = []

async def connect(self):
self.pool = await asyncpg.create_pool(port=int(getenv("POSTGRES_PORT")),
user=getenv("POSTGRES_USER"),
password=getenv("POSTGRES_PASSWORD"),
database=getenv("POSTGRES_DB"))

async def create_schema(self):
async with self.pool.acquire() as connection:
await connection.execute("""
DROP TABLE IF EXISTS table1;

CREATE TABLE table1 (key INTEGER,
value INTEGER,
PRIMARY KEY (key));

CREATE OR REPLACE FUNCTION notify_mychannel()
RETURNS TRIGGER AS
$$
DECLARE
payload TEXT;
BEGIN
-- This is a simplified payload. In the actual program,
-- the payload contains a key and an event label
payload := json_build_object('key', NEW.key);
PERFORM pg_notify('mychannel', payload);
RETURN NULL;
END
$$
LANGUAGE plpgsql;

CREATE TRIGGER update_value
AFTER UPDATE OF value ON table1
FOR EACH ROW
EXECUTE PROCEDURE notify_mychannel();

CREATE TRIGGER new_value
AFTER INSERT ON table1
FOR EACH ROW
EXECUTE PROCEDURE notify_mychannel();
""")

async def add_listener(self, channel, callback):
"""
This wrapper function saves the connection in self.listeners, so that it is not garbage-collected.
"""
connection = await self.pool.acquire()
await connection.add_listener(channel, callback)
self.listeners.append({'connection': connection,
'channel': channel,
'callback': callback})

async def upsert(self, key: int, value: int):
async with self.pool.acquire() as connection:
await connection.execute("INSERT INTO table1(key, value)"
" VALUES ($1, $2)"
" ON CONFLICT (key)"
" DO UPDATE SET value = $2;",
key, value)

async def fetch(self, key: int):
async with self.pool.acquire() as connection: # [1]
record = await connection.fetch("SELECT key, value FROM table1"
" WHERE key = $1;") # [2]
return record['value']

class ExampleApplication:
def __init__(self):
self.database_interface = DatabaseInterface()

async def start(self):
await self.database_interface.connect()
await self.database_interface.create_schema()
await self.database_interface.add_listener('mychannel', self.callback)

async def callback(self, connection: asyncpg.Connection,
pid: int, channel: str, payload: str):
"""
This is an example callback function. In the actual app, the payload contains a key and event value,
which, depending on the event, should trigger fetches from multiple different tables in the database,
return data which will be used in other async operations.
"""
payload_json = json.loads(payload)
value = await self.database_interface.fetch(payload_json['key']) # This is the line that does not complete
print(f"Value {value}") # This line is never run

async def main():
app = ExampleApplication()
await app.start()
await app.database_interface.upsert(1, 10)

if __name__ == "__main__":
asyncio.run(main())
```

コントリビューションガイド

このリポジトリのコントリビューションガイドは索引されていません

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

調査の方向性

まず、提供されている Python 再現コードを実行し、ExampleApplication.callback から DatabaseInterface.fetch までを追跡します。特に [1] でのプール取得と [2] での fetch を確認してください。add_listener が保持する listener 接続と、fetch が取得する接続を比較し、asyncpg の listener と pool の動作を調査します。callback がデータベースの fetch を完了するか、明確な workaround とともに制限事項が文書化されれば完了です。

索引モデルが issue の本文から書いたものです。

評価

技術スタック
postgresql, python
領域
backend, databases
issue の種類
バグ
難易度
4/5
見積もり時間
3〜5日
活発さ
停滞
明瞭さ
おおむね明確
初心者へのやさしさ
28/100

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。