Задачи сельдерея с помощью psycopg: ProgrammingError последняя операция не дала результатаPython

Программы на Python
Anonymous
Задачи сельдерея с помощью psycopg: ProgrammingError последняя операция не дала результата

Сообщение Anonymous »

Я работаю над проектом, в котором у меня есть
  • база данных PostgreSQL 16.2.
  • Бэкэнд Python 3.12. использование psycopg 3.2.1 и psycopg_pool 3.2.2.
  • Celery для обработки асинхронных задач.
Задачи 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()
Код работает в большинстве случаев, но иногда, когда запускается несколько задач одного типа, я получаю ошибки типа psycopg.ProgrammingError: последняя операция не выполнена. t выдать результат при втором вызове fetchone().
Тем временем в базе данных я вижу следующее предупреждение

Код: Выделить всё

WARNING:  there is already a transaction in progress
Я подозреваю, что могут быть какие-то проблемы с тем, как я работаю с соединениями, но не могу их найти.
Как насколько я знаю, после вызова get_db_conn() это соединение недоступно для других задач, поэтому теоретически не может быть несколько задач, использующих одно и то же соединение, и поэтому при выполнении второго fetchone() не должно быть уже выполняемых транзакций. вызов.
Ресурс существует, поскольку любая другая задача может получить к нему доступ, так что проблема не в этом.

Подробнее здесь: https://stackoverflow.com/questions/789 ... -produce-a

Вернуться в «Python»