Pymongo 사용기


Pymongo 사용기

안녕하세요. 실시간 결제 거래 매출 데이터를 관리하는 업무를 하고 있는 주니어 개발자 입니다.

파이썬 환경에서 MongoDB를 어떻게 사용 하며 개인적으로 사용 했던 방식과 레플리카 셋팅이 되지 않은 단일 인스턴스도 트랜잭션에 대한 Rollback을 구현 했던 방식을 공유합니다.

Table of Contents

  1. Python에서 어떻게 MongoDB를 사용하는가
  2. 개인적으로 Pymongo를 다루는 방법
  3. 단일 인스턴스에서 보상 트랜잭션 구현하기
  4. 알아두면 좋을 점

파이썬에서 몽고DB를 다루는 방법

Setup:

mongo // mongo DB 접속

show dbs; // 현재 존재 하는 데이터베이스 목록 확인

use test; // 테스트 데이터베이스가 없더라도 'use'를 이용하여 생성 가능

db.test_collection.insertOne({"test_document": "hello mongo"}); // 데이터 저장

Python에서 MongoDB를 사용하기 위한 코드는 아래와 같습니다.

"""database.py"""

from pymongo import MongoClient

DB_NAME = "test"
COLLECTION = "test_collection"

mongo_client = MongoClient()  # NOTE: 아무런 값을 입력하지 않을 시 localhost가 기본 값으로 설정 됨

mongo_database = mongo_client[DB_NAME]

mongo_collection = mongo_database[COLLECTION]

result = list(mongo_collection.find())  # NOTE: 몽고는 `findOne`이 아닌 여러 데이터를 조회 하면 Cursor 객체를 반환 하기 때문에 List로 Unpacking이 필요함

print(result)  # NOTE: 아무런 값을 입력하지 않을 시 'SELECT * FROM test_collection' 으로 질의 함

""" 조회 결과 """
# [{'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}]

이렇게 간편하게 사용 할 수 있는 장점이 있지만 보다시피 클라이언트 객체, 데이터 베이스 객체, 컬렉션 객체 총 3번을 걸쳐 실제 컬렉션에 존재 하는 도큐먼트에게 DML, DDL이 가능한 형태입니다.

Module 단위로 분리 되어있는 실제 파이썬 환경에서 데이터를 조회 하기 위해 항상 몽고 컬렉션 또는 몽고 데이터베이스 객체를 참조 시키는 행위가 발생하게 됩니다. 예를 들면 아래 처럼 말이죠.

"""module_a.py"""
from database import mongo_database

def module_a_mongo_access(collection_name):
    mongo_collection = mongo_database[collection_name]
    return list(mongo_collection.find({"test_document": "hello mongo"}))

if __name__ == '__main__':
    print(module_a_mongo_access("test_collection"))
 
""" 조회 결과 """
# [{'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}]

이런 코드는 앞으로 조회 하고자 하는 컬렉션이 변경 될 때 마다 데이터베이스 객체에서 컬렉션 이름을 기준으로 컬렉션 객체를 만들어야 합니다.

불편함이 한 두가지가 아닙니다. 만약 서비스 계층에서 해당 레포지토리 객체를 받아서 질의 하는 메서드를 만드는 데 조회 하고자 하는 컬렉션이 다를 때 마다 데이터베이스 객체에서 컬렉션을 뽑아와야 하는 행위가 메서드 단위 마다 생성이 되는게 굉장히 불편하죠.

그 점을 보완 하기 위해서 파이썬 진영 백엔드 프레임워크에서 잘 알려진 장고 프레임워크의 ORM을 살펴보면 db_session과 object_manager를 통해 질의 해야 할 table을 결정하게 됩니다.

이 방법을 따라 하기란 매우 어렵지만 원하는 선 에서의 편의성을 제공하는 정도는 가능하게 될 것입니다.

개인적으로 몽고DB를 다루는 방법

1️⃣ MongoClientManager를 통해 객체 관리하기

MongoDB를 사용하기 위한 객체를 관리하는 클래스로서 역할을 하고 있는데 예제 케이스에서는 MongoDB의 경우 여러 데이터베이스를 옮겨 다니기 보다 하나의 데이터베이스에서 여러 개의 컬렉션을 참조 하는 경우가 많기 때문에 환경 변수나 데이터베이스에 접근 해야 하는 값들을 의존성 주입을 통해 손 쉽게 유연한 구조로 확장 할 수 있습니다.

import logging
from typing import Dict, Optional

from pymongo import MongoClient
from pymongo.database import Database
from pymongo.errors import PyMongoError

from commons.exceptions import EnvException

class MongoClientManager:
    _config = None
    _client: Optional[MongoClient] = None
    _database: Optional[Database] = None

    def __init__(self, target: str, db_config: Dict[str, str]):
        """
        몽고DB 데이터베이스 커넥션 객체 관리

        :param target: AWS 접속 서버 이름
        :param db_config: 데이터베이스 설정 값
        """
        self.target = target
        self._config = db_config

    def connect(self) -> Database:
        """
        Database 생성
        :return: MongoDatabase 객체 반환
        """
        database_url: str = self._config.get("database_url")

        if not database_url:
            raise EnvException("MongoDB URL does not exist")

        # NOTE: 커넥션 연결 대기 시간 10초로 설정
        mongo_client: MongoClient = MongoClient(database_url, serverSelectionTimeoutMS=10000)
        MongoClientManager._client = mongo_client

        database_name: str = self._config.get("database_name")

        if not database_name:
            raise EnvException("Database name does not exist")

        MongoClientManager._database = mongo_client[database_name]

        return self._database

2️⃣ MongoDB Session을 객체의 local 또는 Instance 범위의 변수로 할당하기

Repository 계층에서 db_session을 사용하기 위해 필요한 데코레이터인데요. 아무래도 재사용성을 고려 하다보면 그 대상이 함수일 때, 클래스 일 때 전부 같은 동작을 해야 한다고 판단 했기 때문에 데코레이터 함수 내에서 전부 커버하고 있는 모습을 나타냅니다.

예제는 굉장히 러프하게 코드가 구현 되어 있지만 견고하게 하기 위해서는 필수로 타입 체크 또는 생각되는 다른 예외 케이스들에 대한 검증을 많이 거치는 것이 필요합니다.

def use_db(obj):
    """
    데이터베이스 객체를 대상 함수 또는 클래스에게 주입 하는 데코레이터
    :param obj: 데이터베이스를 사용 해야 하는 함수 또는 클래스
    - 함수: 인자에 `db_session=None` 을 포함
    - 클래스: 클래스 변수에 `db_session=None` 을 포함
    """
    if isinstance(obj, type):  # NOTE: 데코레이터 대상이 클래스 객체 일 때
        original_init = obj.__init__

        def class_decorator(instance, *args, **kwargs):
            db = MongoClientManager.get_database()
            instance.db_session = db
            original_init(instance, *args, **kwargs)

        obj.__init__ = class_decorator
        return obj
    else:

        def function_decorator(*args, **kwargs):  # NOTE: 데코레이터 대상이 함수 객체 일 때
            db = MongoClientManager.get_database()
            result = obj(db_session=db, *args, **kwargs)
            return result

        return function_decorator

3️⃣ 프로젝트 환경에서 사용 할 Query 추상화 객체 만들기

메서드 내용은 대부분 사용하는 코드들이기 때문에 크게 반감 없이 읽으실 수 있으실 겁니다.

이렇게 Pymongo 코드를 래핑 했을 때 로깅을 견고하게 할 수 있고 데이터 검증을 내부에서 수행 하여 서비스 계층에서 별도의 Validation 없이 데이터를 사용 할 수 있는 장점이 있습니다.

import logging

from bson import ObjectId

from commons.custom import DataDict
from commons.decorators import use_db
from commons.exceptions import DatabaseCreateError, DatabaseUpdateError, EmptyQueryResult

logger = logging.getLogger(__name__)

class CommandManager:
    def __init__(self, collection):
        self.collection = collection

    def get(self, **kwargs):
        """
        쿼리 결과가 단일 일 경우 조회 메서드
        """

        objects = None

        try:
            objects = self.collection.find_one(kwargs)
        except Exception as e:
            logger.warning(f"command `get` method during occured error, detail: {e}")

        return objects

    def all(self, **kwargs):
        """
        쿼리 결과가 다중일 경우 조회 메서드
        """

        objects = None
        try:
            objects = list(self.collection.find(kwargs))

            if len(objects) == 0:
                return None

        except Exception as e:
            logger.warning(f"command `all` method during occured error, detail: {e}")

        return objects

    def create(self, **kwargs):
        """
        데이터 생성 메서드
        """
        try:
            object_id = self.collection.insert_one(kwargs)
            created = True
        except Exception as e:
            raise DatabaseCreateError(f"command `create` method during occured error, detail: {e}")
        return object_id, created

    def update(self, update_content, **kwargs):
        """
        기존에 존재하는 데이터 업데이트 메서드
        """
        try:
            matched_count = self.all(**kwargs) or []
            if len(matched_count) != 1:
                raise EmptyQueryResult(f"match query: {kwargs} has {len(matched_count)} of selected.")

            return self.collection.update_one(kwargs, {"$set": update_content}, upsert=False)

        except Exception as e:
            raise DatabaseUpdateError(f"command `update` method during occured error, detail: {e}")

4️⃣ 쿼리 매니저와 결합 할 Repository 구성하기

아래와 같은 객체를 구성 함으로써 서비스 객체에서 해당 Repository를 주입 하거나, 합성을 통해 접근이 가능하게 됩니다.

@use_db
class Repository:

    db_session = None

    def __getattr__(self, collection_name):
        return CommandManager(self.db_session[collection_name])
MongoClientManager.global_connect()

if __name__ == '__main__':
    mongo_repository = Repository()
    print(mongo_repository.test_collection.get(test_document="hello mongo"))

# 만약 test_collection 이 아닌 다른 컬렉션에서 데이터를 조회 하고 싶을 때 컬렉션 명만 수정하면 됩니다.
# mongo_repository.another_collection.get(query) 이런식으로 사용 할 수 있고 query는 Dictionary 입니다.

"""조회 결과"""
# {'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}

단일 인스턴스에서 보상 트랜잭션 구현하기

Mongo DB는 Replication 환경에서 논리적 Rollback을 지원합니다. 사내 환경은 운영은 Replicaset으로 구성 되어 있지만 QA 서버는 단일 인스턴스에서 동작합니다.

그렇기 때문에 QA 서버에서 정상 동작을 하지 않는 함수를 만드는건 옳지 않다고 판단 하여, 단일 인스턴스에서 논리적 Rollback을 구현 하는 방법을 소개 하고자 합니다.

class CommandManager:
    def __init__(self, collection):
        self.collection = collection

    def update(self, update_content, **kwargs):
            """
            기존에 존재하는 데이터 업데이트 메서드
            """
            try:
                matched_count = self.all(**kwargs) or []
                if len(matched_count) != 1:
                    raise EmptyQueryResult(f"match query: {kwargs} has {len(matched_count)} of selected.")
    
                return self.collection.update_one(kwargs, {"$set": update_content}, upsert=False)

        except Exception as e:
            raise DatabaseUpdateError(f"command `update` method during occured error, detail: {e}")

    def rollback(self, before_update):
        """
        업데이트 했던 내용 되돌리기
        :param before_update: 업데이트 했던 기존 데이터
        :type before_update: dict[str, dict]
        :return: Rollback이 수행 된 대상 ObjectId
        :rtype: list[str]
        """
        rollback_id_list = []

        for object_id, update_content in before_update.items():
            filter_query = {"_id": ObjectId(object_id)}
            self.update(update_content, **filter_query)
            rollback_id_list.append(object_id)

        return rollback_id_list

if __name__ == "__main__":
    mongo_repository = Repository()
    rollback_id_list = list(map(mongo_repository.test_collection.rollback, before_update_list))

Rollback을 하기 위한 메서드를 생성 하고, 업데이트 이 전 데이터를 {Mongo ObjectId: 업데이트 이전 데이터(딕셔너리)} 상태로 갖고 있다면 얼마든지 기존 데이터로 돌아갈 수 있죠.

예를 들면 아래와 같습니다.

before_update = {ObjectId("1234567890"): {"test": "테스트"}}
after_update = {ObjectId("1234567890"): {"test": "테스트2"}}

for oid_, content in before_update.items():
    query = {"_id": oid}  # NOTE: "1234567890"

    """
    UPDATE test_collection
    SET test = "테스트"
    WHERE _id = "1234567890"
    """
    self.update(content, **query)

어떤 데이터를 Rollback 할 것인가를 정의 하여 위 처럼 업데이트 이전 데이터를 메모리에 저장 하고 있다면 언제든 되돌릴 수 있게 됩니다.

그 외 Pymongo를 사용 하면 알면 좋은 정보

💡 Pymongo의 Cursor의 Lazy evaluation 방식

“which can iterate over the results a single time” — [StackOverflow]

print(result)
print(list(result))
print(result.retrieved)
print(list(result))

""" 
조회 결과

<pymongo.cursor.Cursor object at 0x103e27b90>
[{'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}]
1
[]
""""

find메서드는 Cursor객체를 반환하며 Lazy evaluate 방식을 사용 하고 있습니다.

그렇기 때문에 한 번 조회 된 커서 객체의 데이터는 다시 조회 하면 비어있는 것을 알 수 있는데, 이걸 방지 하고 싶다면 명시적으로 rewind 메서드로 Unevaluated 상태로 변경 하거나 메모리에 올려놓고 사용하는 것을 권장합니다.

print("first evaluate")
print(result)
print(list(result))
print(result.retrieved)

print("second evaluate")
result.rewind()
print(result)
print(result.retrieved)
print(list(result))

"""
조회 결과

first evaluate
<pymongo.cursor.Cursor object at 0x105df7b60>
[{'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}]
1

second evaluate
<pymongo.cursor.Cursor object at 0x105df7b60>
0
[{'_id': ObjectId('66d2c0f80b6c06da3d815ed9'), 'test_document': 'hello mongo'}]
"""

그렇기 때문에 대부분 Cursor 객체의 Unpacking 과정을 list로 푼 뒤 변수에 담아두고 사용합니다.

만약 데이터가 너무 많다면 next를 이용해서 가공 하되 재사용이 불가능 한 것을 인지하고 있어야 합니다.

💡 Pymongo `InsertOne` Bson Object Size limit 16mb

실제 프로젝트에서 배치 프로세스를 통해 대용량 데이터를 Python의 Dictionary와 내부 속성의 데이터 리스트를 저장 하기 위해 InsertOne 명령어를 날리면 DocumentTooLarge 에러가 발생하게 됩니다. 이는 MongoDB에서 의도한 데이터 크기이며 만일 해당 에러가 발생 했다면 데이터를 청크로 쪼개서 넣거나 GridFS API를 이용하라고 안내 하고 있습니다.

한 때 아래 코드 처럼 청크로 쪼개서 사용 했었는데 좋은 방법은 아닌 것 같습니다.

그렇기 때문에 개인적으로 하나의 도큐먼트가 16MB가 넘어야 하는 이유와 해당 컬렉션에 대한 설계를 다시 해 보는 것을 고려하고 있습니다.

class MongoUtil:
    LIMIT_DOC_SIZE: int = 300  # NOTE: 청크 데이터 사이즈 확인 후 갯수 조정

    @staticmethod
    def is_data_size_limit(data: List[Dict]) -> bool:
        """
        몽고DB에 저장 할 데이터의 크기가 16MB가 초과 하는지 여부 검사
        :param data: 몽고DB에 저장 할 데이터
        :return: bool -> False일 시 데이터 삽입 가능
        """

        return len(data) > MongoUtil.LIMIT_DOC_SIZE

    @staticmethod
    def split_data(data: List[Dict]):
        Data = namedtuple("Data", "current, nxt")
        current = data[: MongoUtil.LIMIT_DOC_SIZE]
        nxt = data[MongoUtil.LIMIT_DOC_SIZE:]

        return Data(current, nxt)

이 내용과 관련하여 읽어보면 좋은 글들