event loop for database engine
This commit is contained in:
@@ -68,12 +68,14 @@ def _run_process_loop_task(task_id: str, runner: LoopRunner) -> None:
|
||||
|
||||
|
||||
async def _run_process_loop_task_async(task_id: str, runner: LoopRunner) -> None:
|
||||
from database.db import reset_async_db_engine
|
||||
from bot import bot
|
||||
from core.bootstrap import bootstrap
|
||||
from database import async_session_maker, init_db
|
||||
|
||||
logger.info("[PeriodicManager] Process-loop задача {} запущена, PID={}", task_id, os.getpid())
|
||||
try:
|
||||
reset_async_db_engine()
|
||||
await init_db()
|
||||
await bootstrap()
|
||||
await runner(bot, async_session_maker)
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
from .audit import *
|
||||
from .bans import *
|
||||
from .coupons import *
|
||||
from .db import async_session_maker
|
||||
from .db import Base, async_session_maker, engine, reset_async_db_engine
|
||||
from .gifts import *
|
||||
from . import identities
|
||||
from .hot_leads import *
|
||||
|
||||
+23
-11
@@ -19,17 +19,23 @@ if USE_PGBOUNCER and "+asyncpg" in DATABASE_URL:
|
||||
_db_url = f"{_db_url}{sep}prepared_statement_cache_size=0"
|
||||
|
||||
_pool_recycle = 60 if USE_PGBOUNCER else 300
|
||||
engine = create_async_engine(
|
||||
_db_url,
|
||||
echo=False,
|
||||
future=True,
|
||||
pool_size=DB_POOL_SIZE,
|
||||
max_overflow=DB_MAX_OVERFLOW,
|
||||
pool_timeout=60,
|
||||
pool_pre_ping=True,
|
||||
pool_recycle=_pool_recycle,
|
||||
connect_args=_connect_args,
|
||||
)
|
||||
|
||||
|
||||
def _create_engine():
|
||||
return create_async_engine(
|
||||
_db_url,
|
||||
echo=False,
|
||||
future=True,
|
||||
pool_size=DB_POOL_SIZE,
|
||||
max_overflow=DB_MAX_OVERFLOW,
|
||||
pool_timeout=60,
|
||||
pool_pre_ping=True,
|
||||
pool_recycle=_pool_recycle,
|
||||
connect_args=_connect_args,
|
||||
)
|
||||
|
||||
|
||||
engine = _create_engine()
|
||||
|
||||
async_session_maker = async_sessionmaker(
|
||||
bind=engine,
|
||||
@@ -42,6 +48,12 @@ Base = declarative_base()
|
||||
WARM_POOL_COUNT = 10
|
||||
|
||||
|
||||
def reset_async_db_engine() -> None:
|
||||
global engine
|
||||
engine = _create_engine()
|
||||
async_session_maker.configure(bind=engine)
|
||||
|
||||
|
||||
async def warm_pool() -> None:
|
||||
"""
|
||||
Прогревает пул соединений при старте.
|
||||
|
||||
+3
-3
@@ -3,15 +3,15 @@ from datetime import datetime
|
||||
from sqlalchemy import select
|
||||
|
||||
from config import ADMIN_ID
|
||||
from database.db import async_session_maker, engine
|
||||
from database import db
|
||||
from database.models import Admin, Base, User
|
||||
|
||||
|
||||
async def init_db():
|
||||
async with engine.begin() as conn:
|
||||
async with db.engine.begin() as conn:
|
||||
await conn.run_sync(Base.metadata.create_all)
|
||||
|
||||
async with async_session_maker() as session:
|
||||
async with db.async_session_maker() as session:
|
||||
result = await session.execute(select(User).where(User.tg_id == 0))
|
||||
if not result.scalar_one_or_none():
|
||||
session.add(
|
||||
|
||||
Reference in New Issue
Block a user