[WAFFLE] 04-3 SQLAlchemy
Published:
In this post, concepts of SQLAlchemy will be provided.
SQLAlchemy
과제에서 사용한 io는 db와 소통한 SQLAlchemy 밖에 없다. SQLAlchemy 내부에는 비동기 API도 있는데 이를 사용해야 제대로 FastAPI를 활용할 수 있다.
from typing import Any, Generator
from sqlalchemy import create_engine
from sqlalchemy.orm import Session, sessionmaker
from database.coomon import Base
class DatabaseManager:
def __init__(self):
# TODO pool 이 뭘까요?
# pool_recycle 은 뭐고 왜 28000으로 설정해두었을까요?
self.engine = create_engine(
DB_SETTINGS.url,
pool_recycle=28000,
pool_pre_ping=True,
)
self.session_factory = sessionmaker(bind=self.engine, expire_on_commit=False)
def get_db_session() -> Generator[Session, Any, None]:
session = DatabaseManager().session_factory()
try:
yield session
session.commit()
except Exception as e:
session.rollback()
raise e
finally:
session.close()
SQLAlchemy는 가장 로우 레벨에서는 engine이 있고 connection이 있음 engine은 어떤 db와 연결하기 위한 연결정보를 저장하는 곳. 그 연결정보를 바탕으로 connection을 만들어낼 수 있는 것. 여러개의 db와 연결하려면 여러 개의 engine 사용. engine을 만들고 그 engine을 이용해서 sessionmaker를 만듦. session_factory에 이를 저장해두고 get_db_session에서 활용.
class UserStore:
def __init__(self, session: Annotated[Session, Depends(get_db_session)]) -> None:
self.session = session
여기서 sesion을 종속성으로 주입받음. get_db_session을 여기서 호출됨.
종속성으로 session이 필요할 때 get_db_session을 호출해서 session_factory로 session을 만들고 그걸 yield를 하고 요청에 대한 응답이 반환이 되면 그 뒤에 코드가 실행이 되면서 롤백이나 세션 클로스가 실행됨. 종속성 방식의 장점은 하나의 세션을 여러 개 요청이 공유하게 되면 동시성 문제가 발생함. 각각의 요청은 고유한 세션을 받아야함. 그런 관점에서 종속성으로 세션을 관리하면 알아서 매번 세션을 만들어주고 관리해주므로 문제 없음. 그러나 동기 세션이기 때문에 비동기 함수 내부에서 호출하게 되면 이벤트 루프 전체가 블락 되는 상황이 발생함. 그래서 이부분을 비동기 api로 바꿔야 함.
from typing import Any, Generator
from sqlalchemy import create_engine
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy.ext.asynicio import AsyncSession, async_session_maker, create_async_engine
from database.coomon import Base
class DatabaseManager:
def __init__(self):
# TODO pool 이 뭘까요?
# pool_recycle 은 뭐고 왜 28000으로 설정해두었을까요?
self.engine = create_async_engine(
DB_SETTINGS.url,
pool_recycle=28000,
pool_pre_ping=True,
)
self.session_factory = async_sessionmaker(bind=self.engine, expire_on_commit=False)
async def get_db_session() -> AsyncGenerator[AsyncSession, Any, None]:
session = DatabaseManager().session_factory()
try:
yield session
await session.commit()
except Exception as e:
await session.rollback()
raise e
finally:
await session.close()
- from sqlalchemy.ext.asynicio import AsyncSession, async_session_maker, create_async_engine 추가.
- 일반 engine을 async로 바꿈, sesion_maker도 async로 바꿈.
- get_db_session도 수정.
- Store, service, views 도 변경.
class UserStore: def __init__(self, session: Annotated[AsyncSession, Depends(get_db_session)]) -> None: self.session = session async def get_user_by_username(self, username: str) -> User | None: return await self.session.scalar(select(User).where(User.username == username)) def update_user( self, username: str, email: str | None, address: str | None, phone_number: str | None, ) -> User: user = await self.get_user_by_username(username)코루틴 함수로 정의되면 비동기 함수이고 비동기 함수는 실제로 값을 반환하지 않고 코루틴 객체를 반환하므로 await를 이용해서 코루틴 객체를 실행하고, 그 결과값을 가져올 수 있도록 해야함. 코드 마지막 줄에서 await를 적어야 하는 이유.
이렇게 함으로써 단일 코어에서 실행되는 서버라도 여러 개 요청을 동시에 처리할 수 있는 최소환의 기능을 갖추게 되었다.
하나의 api 요청에서 db와 여러 번 통신할 일이 있다. 원하는 만큼 묶어서 실행하고 싶은데 그 묶는 단위를 transaction이라 함. 현재 코드는 마지막에만(요청이 종료된 직후에만) commit 하도록 되어 있는데 우리가 원할 때 commit 할 수 있도록 코드를 바꾸어보자.
it = get_db_session()
session = await next(it)
get_db_session은 generator이기 때문에 우리가 low 하게 사용한다 하면 FastAPI 내부적으로 위와 같이 session을 받았을 것이다.
# Store
async def add_user(self, username: str, password: str, email: str):
with transaction() as session2:
if await self.get_user_by_username(username):
raise UsernameAlreadyExistsError()
..
# connection
@asynccontextmanager
async def transaction() -> AsyncGenerator[AsyncSession, None]
session = DatabaseManager().session_factory()
try:
yield session
await session.commit()
except Exception as e:
await session.rollback()
raise e
finally:
await session.close()
위와 같이 적으면 우리가 명시적으로 어디부터 어디까지 session2가 살아있을지를 나타낼 수 있는데 추가로 몇가지 변형이 코드에 더 필요하다. 하지만 이 방식은 다소 코드를 읽기 불편하기 때문에 다른 방식을 알아보자.
# store
class UserStore:
async def add_user(...):
async with transaction():
...
# connection
session = async_scoped_session(session_factory=DatabaseManager().session_factory, scopefunc=asyncio.current_task)
@asynccontextmanager
async def transaction() -> AsyncGenerator[None, None]
try:
yield
await session.commit()
except Exception as e:
await session.rollback()
raise e
finally:
...
기본적으로 우리가 사용하는 방식은 session_factory에서 session을 만들어서 넘겨주고 이를 이용해서 insert 등을 하는 것. 위 방식은 session을 만들어서 건네주는 것이 아니라 전역 변수 session을 사용하는 것. 그런데 아까 요청들끼리 동일한 session을 공유하면 동시성 문제가 발생한다고 함. 그래서 get_db_session 같은 종속성 함수를 만들어서 요청이 들어올 때마다 session을 새로 생성하고 요청이 끝나면 session이 자동으로 닫히게 사용했던 것. 그러나 async_scoped_session이 반환하는 것은 함수가 아니라 클래스이다. 내부적으로 파보면 프록시라는 함수를 호출해 문제 없이 동작하게 함. 그 과정에서 핵심이 scopefunc이 요청마다 다른 값을 보내는 함수여야 함. 그게 key 값으로 쓰이기 때문에. asyncio.current_task는 현재 실행되고 있는 task의 이름을 반환함.
지금까지 한 것을 정리하면 일단 우리가 구현한 함수를 전부 async await으로 바꿈. 그걸 하기 앞서 SQLAlchemy에서 session, engine을 async session, async engine으로 바꿈. 그리고 종속성에서 반환하는 값이 async session이니까 주입받는 값도 전부 async session으로 바꾸고 async session에 있는 비동기 API를 활용하기 위해서 store.py 함수들을 전부 async으로 바꿔주고, SQLAlchemy API호출할때 await 붙여주고 service.py에서도 store 함수 호출할 때 await 붙여주고 async로 바꿔주고 views.py에서도 async 함수로 바꿔주고 service 호출하는 부분 await로 바꿔줌. 이번에는 get_db_session은 요청 단위로만 transaction이 가능하다보니 우리가 원하는 단위로 바꾸기 위해 transactioncontextmanger를 만듦.

Leave a Comment