пул соединений psycopg2 (min=1 max=3), db.put_conn вместо close, v1.20

This commit is contained in:
2026-06-17 18:23:41 +04:00
parent dfb9a4218b
commit 2615e781ce
3 changed files with 51 additions and 35 deletions
+46 -30
View File
@@ -1,27 +1,38 @@
"""
db.py — Транспортный слой к базе данных.
db.py — Транспортный слой к базе данных (с пулом соединений).
Только connect() и query(). Никакой бизнес-логики, никакого DDL.
DDL → schema.py. Бизнес-запросы → upload.py, extractor.py, api.py.
Функции:
_pg_connect(dbname) — сырое подключение к ЛЮБОЙ базе (используется /test createdb)
connect() — подключение к ЦЕЛЕВОЙ базе (из DB_NAME)
query(sql, params) — выполнить SELECT → (result, error)
execute(sql, params) — выполнить INSERT/UPDATE/DELETE → (rowcount, error)
query_one(sql, params) — как query, но возвращает одну строку или None
Пул psycopg2.pool.ThreadedConnectionPool (min=1, max=3).
Все query/execute забирают соединение из пула и возвращают обратно.
"""
import os
from psycopg2 import pool as _pgpool
import psycopg2
_pool = None
def _get_pool():
"""Ленивая инициализация пула."""
global _pool
if _pool is None:
_pool = _pgpool.ThreadedConnectionPool(
minconn=1,
maxconn=3,
host=os.getenv("DB_HOST"),
port=os.getenv("DB_PORT", 5432),
dbname=os.getenv("DB_NAME"),
user=os.getenv("DB_USER"),
password=os.getenv("DB_PASS"),
sslmode=os.getenv("DB_SSLMODE", "disable"),
)
return _pool
def _pg_connect(dbname):
"""
Сырое подключение к указанной базе данных.
Используется:
- connect() для целевой БД
- test_routes.py для /test createdb (подключение к 'postgres')
Сырое подключение к ЛЮБОЙ базе (для /test createdb).
НЕ из пула — для создания БД нужна отдельная сессия.
"""
return psycopg2.connect(
host=os.getenv("DB_HOST"),
@@ -35,20 +46,26 @@ def _pg_connect(dbname):
def connect():
"""
Подключение к ЦЕЛЕВОЙ базе данных (DB_NAME из переменных окружения).
Возвращает (connection, None) при успехе или (None, error) при ошибке.
Подключение к ЦЕЛЕВОЙ базе из пула.
Возвращает (connection, None) или (None, error).
"""
try:
return _pg_connect(os.getenv("DB_NAME")), None
return _get_pool().getconn(), None
except Exception as e:
return None, str(e)
def put_conn(conn):
"""Вернуть соединение в пул."""
try:
_get_pool().putconn(conn)
except Exception:
pass
def query(sql_text, params=None):
"""
Выполнить SQL-запрос к целевой БД.
Возвращает (result, None) или (None, error).
SELECT → (result, None) или (None, error).
result = {"columns": [...], "rows": [[...], ...]}
"""
conn, err = connect()
@@ -60,17 +77,16 @@ def query(sql_text, params=None):
rows = cur.fetchall()
cols = [desc[0] for desc in cur.description] if cur.description else []
cur.close()
conn.close()
return {"columns": cols, "rows": [list(r) for r in rows]}, None
except Exception as e:
conn.close()
return None, str(e)
finally:
put_conn(conn)
def execute(sql_text, params=None):
"""
Выполнить INSERT/UPDATE/DELETE к целевой БД.
Возвращает (rowcount, None) или (None, error).
INSERT/UPDATE/DELETE → (rowcount, None) или (None, error).
"""
conn, err = connect()
if err:
@@ -81,23 +97,23 @@ def execute(sql_text, params=None):
conn.commit()
rowcount = cur.rowcount
cur.close()
conn.close()
return rowcount, None
except Exception as e:
conn.rollback()
conn.close()
try: conn.rollback()
except: pass
return None, str(e)
finally:
put_conn(conn)
def query_one(sql_text, params=None):
"""
Выполнить SELECT и вернуть ПЕРВУЮ строку или None.
Возвращает (row_dict, None) или (None, error).
SELECT одной строки → (row_dict, None) или (None, error).
"""
result, err = query(sql_text, params)
if err:
return None, err
rows = result["rows"]
if not rows:
return None, None # нет строк — не ошибка
return None, None
return dict(zip(result["columns"], rows[0])), None