- база данных PostgreSQL 16.2.
- Бэкэнд Python 3.12. использование psycopg 3.2.1 и psycopg_pool 3.2.2.
- Celery для обработки асинхронных задач.
Код: Выделить всё
import os
from psycopg_pool import ConnectionPool
from contextlib import contextmanager
PG_USERNAME = os.getenv('PG_USERNAME')
if not PG_USERNAME:
raise ValueError(f"Invalid postgres username")
PG_PASSWORD = os.getenv('PG_PASSWORD')
if not PG_PASSWORD:
raise ValueError(f"Invalid postgres pass")
PG_HOST = os.getenv('PG_HOST')
if not PG_HOST:
raise ValueError(f"Invalid postgres host")
PG_PORT = os.getenv('PG_PORT')
if not PG_PORT:
raise ValueError(f"Invalid postgres port")
# Options used to prevent closed connections
# conn_options = f"-c statement_timeout=1800000 -c tcp_keepalives_idle=30 -c tcp_keepalives_interval=30"
conninfo = f'host={PG_HOST} port={PG_PORT} dbname=postgres user={PG_USERNAME} password={PG_PASSWORD}'
connection_pool = ConnectionPool(
min_size=4,
max_size=100,
conninfo=conninfo,
check=ConnectionPool.check_connection,
#options=conn_options,
)
@contextmanager
def get_db_conn():
conn = connection_pool.getconn()
try:
yield conn
finally:
connection_pool.putconn(conn)
Код: Выделить всё
@app.task(bind=True)
def example_task(self, id):
with get_db_conn() as conn:
try:
with conn.cursor(row_factory=dict_row) as cursor:
test = None
cursor.execute('SELECT * FROM test WHERE id = %s', (id,))
try:
test = cursor.fetchone()
except psycopg.errors.ProgrammingError:
logger.warning(f'Test log msg')
conn.rollback()
return
cursor.execute("UPDATE test SET status = 'running' WHERE id = %s", (id,))
conn.commit()
# Some processing...
# Fetch another resource needed
cursor.execute('SELECT * FROM test WHERE id = %s', (test['resource_id'],))
cursor.fetchone()
# Update the entry with the result
cursor.execute("""
UPDATE test
SET status = 'done', properties = %s
WHERE id = %s
""", (Jsonb(properties), id))
conn.commit()
except Exception as e:
logger.exception(f'Error: {e}')
conn.rollback()
with conn.cursor(row_factory=dict_row) as cursor:
# Update status to error with exception information
cursor.execute("""
UPDATE test
SET status = 'error', error = %s
WHERE id = %s
""", (Jsonb({'error': str(e), 'stacktrace': traceback.format_exc()}), webpage_id))
conn.commit()
Тем временем в базе данных я вижу следующее предупреждение
Код: Выделить всё
WARNING: there is already a transaction in progressКак насколько я знаю, после вызова get_db_conn() это соединение недоступно для других задач, поэтому теоретически не может быть несколько задач, использующих одно и то же соединение, и поэтому при выполнении второго fetchone() не должно быть уже выполняемых транзакций. вызов.
Ресурс существует, поскольку любая другая задача может получить к нему доступ, так что проблема не в этом.
Подробнее здесь: https://stackoverflow.com/questions/789 ... -produce-a