FastAPI·SQLModel·Scrapy·Celery에서 패턴이 맞물리는 방식

Repository+Unit of Work, 의존성 주입, Strategy, Chain of Responsibility, Command+Retry+Idempotency가 Python 백엔드·크롤링 파이프라인에서 실제로 어떻게 조합되는지 코드로 정리한다.

2026-08-13 · 최초 발행 2025-10-14

Python 기반 백엔드·데이터 파이프라인에서는 재사용성과 일관성을 확보하기 위해 디자인 패턴 적용이 필요해진다. 디자인 패턴은 반복되는 설계 문제에 대한 검증된 해결 템플릿으로, 시스템 결합도를 최소화하고 변경 용이성을 극대화하는 것이 목적이다. 이 글은 FastAPI, SQLModel, Scrapy, Celery 조합을 기준으로 API 계층·영속성 계층·크롤링·비동기 작업 사이의 일관된 의존성 역전, 트랜잭션, 오류 복구, 멱등성을 어떻게 구현하는지 코드로 정리한다.

이 스택에서 패턴들이 맞물리는 지점

Repository는 ORM 세부 구현(SQLModel/SQLAlchemy)을 캡슐화하고 도메인 중심 인터페이스를 제공해 테스트 격리와 교체를 쉽게 한다. Unit of Work(UoW)는 트랜잭션 경계를 관리한다 — 성공하면 commit, 실패하면 rollback하며, Celery 작업이나 API 요청 단위의 일관성을 유지한다. 트랜잭션 경계를 명확히 해두면 데이터 불일치 이슈도 현저히 줄어든다.

의존성 주입은 FastAPI의 Depends/yield로 구현된다. UoW·Repository·Service 객체의 수명 주기를 제어해 구성과 사용을 분리하고, 테스트 시에는 목(mock)을 주입해 빠른 단위 테스트를 가능하게 한다. 레포지토리·서비스를 이렇게 공통화하면 API·작업 코드 중복이 30% 이상 감소할 것으로 예상된다.

Strategy는 Scrapy Spider에서 사이트별·도메인별 파싱 규칙을 분리하는 데 쓰인다. 런타임에 전략을 교체할 수 있고, 신규 사이트 확장이 쉬워지며, 공통 인터페이스로 파이프라인·후처리 일관성을 유지한다. 내부 인터페이스를 재사용하는 구조라면 신규 사이트를 추가할 때 파서 확장 리드타임이 20~40% 단축될 것으로 예상된다.

Chain of Responsibility는 Scrapy Item Pipeline을 필터·정규화·검증·전송 단계로 분리하는 데 쓰인다. 각 핸들러가 조건에 따라 처리하거나 다음으로 위임하며, 오류 항목을 격리하고 모니터링 지점을 명확히 한다.

Command + Retry + Idempotency는 Celery Task를 Command로 모델링하는 조합이다. 실패 시 지수 백오프로 재시도하고, 멱등키(SKU 등)와 고유 제약 조건으로 중복 처리를 예방한다. Outbox나 이벤트 발행 시에는 트랜잭션 경계 내 처리 또는 재시도·보상 트랜잭션을 고려한다. 이 조합을 적용하면 중복·유실 건이 고유 제약과 보정 로직을 기준으로 90% 이상 감소할 것으로 기대된다.

데이터가 흐르는 경로

enqueueUoWQueryEvent/Outboxerror/retryDIScrapy Spider(Strategy)Item Pipelines(Chain of Responsibility)Celery BrokerCelery Worker(Command+Retry)(DB:SQLModel/SQLAlchemy)FastAPI(Read API)(Message/Event Bus)Service Layer(Repository)

입력은 Spider가 수집한 데이터와 API 요청이다. 처리는 파이프라인 정규화 → 큐 적재 → 워크플로우 실행(UoW/Repo) 순으로 진행되고, 출력은 DB 반영, 이벤트 발행, API 응답이다. 일시적 오류는 Celery 재시도·백오프로, 영속 오류는 Dead-letter 큐나 에러 테이블로 격리한다. 트랜잭션은 UoW 경계 안에서 upsert하고, 고유 제약이 충돌하면 보정 업데이트하며, 가능한 한 짧게 유지한다.

코드로 본 구현

전제조건은 Python 3.11+, FastAPI, Uvicorn, SQLModel, SQLAlchemy, Scrapy, Celery, Redis 또는 RabbitMQ(브로커), SQLite/PostgreSQL(데이터베이스)이며, 버전은 예시이므로 최신 정보를 확인해야 한다.

설치는 pip install fastapi uvicorn sqlmodel sqlalchemy scrapy celery[redis] pydantic로 하고, 개발·테스트에는 SQLite/Redis를, 운영에는 PostgreSQL/RabbitMQ를 권장한다.

도메인 모델과 Repository/UoW

# models.py
from typing import Optional
from sqlmodel import SQLModel, Field

class Product(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    sku: str = Field(index=True, unique=True)
    name: str
    price: float
# db.py
from sqlmodel import create_engine, Session, SQLModel
DATABASE_URL = "sqlite:///./app.db"  # 운영은 PostgreSQL 권장: postgresql+psycopg://user:pw@host/db
engine = create_engine(DATABASE_URL, echo=False)

def init_db():
    SQLModel.metadata.create_all(engine)
# repository.py
from typing import Optional
from sqlmodel import Session, select
from sqlalchemy.exc import IntegrityError
from models import Product

class ProductRepository:
    def __init__(self, session: Session):
        self.session = session

    def get_by_sku(self, sku: str) -> Optional[Product]:
        return self.session.exec(select(Product).where(Product.sku == sku)).one_or_none()

    def add(self, product: Product) -> Product:
        self.session.add(product)
        return product

    def upsert_by_sku(self, sku: str, name: str, price: float) -> Product:
        obj = self.get_by_sku(sku)
        if obj:
            obj.name, obj.price = name, price
            return obj
        try:
            obj = Product(sku=sku, name=name, price=price)
            self.session.add(obj)
            return obj
        except IntegrityError:
            self.session.rollback()
            # 경쟁 상황에서 생성 충돌 발생 시 재조회 후 갱신
            obj = self.get_by_sku(sku)
            if obj:
                obj.name, obj.price = name, price
                return obj
            raise
# uow.py
from contextlib import AbstractContextManager
from sqlmodel import Session
from db import engine
from repository import ProductRepository

class UnitOfWork(AbstractContextManager):
    def __init__(self):
        self.session: Session | None = None
        self.products: ProductRepository | None = None

    def __enter__(self):
        self.session = Session(engine)
        self.products = ProductRepository(self.session)
        return self

    def commit(self):
        if self.session:
            self.session.commit()

    def rollback(self):
        if self.session:
            self.session.rollback()

    def __exit__(self, exc_type, exc, tb):
        try:
            if exc_type:
                self.rollback()
            else:
                self.commit()
        finally:
            if self.session:
                self.session.close()
# service.py
from uow import UnitOfWork

class ProductService:
    def __init__(self, uow: UnitOfWork):
        self.uow = uow

    def upsert(self, sku: str, name: str, price: float):
        self.uow.products.upsert_by_sku(sku, name, price)
        # 추가 도메인 규칙/이벤트 발행 위치

FastAPI: DI 기반 엔드포인트

# api.py
from fastapi import FastAPI, Depends
from pydantic import BaseModel
from uow import UnitOfWork
from service import ProductService
from db import init_db

app = FastAPI()

class ProductIn(BaseModel):
    sku: str
    name: str
    price: float

def get_uow():
    with UnitOfWork() as uow:
        yield uow  # 성공 시 commit, 예외 시 rollback

@app.on_event("startup")
def on_startup():
    init_db()

@app.post("/products")
def upsert_product(body: ProductIn, uow: UnitOfWork = Depends(get_uow)):
    svc = ProductService(uow)
    svc.upsert(body.sku, body.name, body.price)
    return {"status": "ok", "sku": body.sku}

실행은 uvicorn api:app --reload로 한다.

Celery: Command + Retry + Idempotency

# tasks.py
import os
from celery import Celery
from uow import UnitOfWork
from service import ProductService

broker_url = os.getenv("CELERY_BROKER_URL", "redis://localhost:6379/0")
backend_url = os.getenv("CELERY_RESULT_BACKEND", "redis://localhost:6379/0")

celery_app = Celery("tasks", broker=broker_url, backend=backend_url)

@celery_app.task(bind=True, autoretry_for=(Exception,), retry_backoff=True, retry_jitter=True, max_retries=5)
def upsert_product_task(self, payload: dict):
    # Command 객체에 해당하는 작업 단위
    with UnitOfWork() as uow:
        svc = ProductService(uow)
        # 멱등키: sku. 고유 제약+upsert로 중복 처리 방지
        svc.upsert(payload["sku"], payload["name"], float(payload["price"]))

실행은 celery -A tasks.celery_app worker --loglevel=INFO로 한다.

멱등성은 고유 제약(SKU) + upsert + 재시도로 확보하고, 태스크 타임아웃·하드타임리밋을 설정하는 것이 권장된다. 재시도 한계를 초과하면 DLQ(Dead-letter Queue)로 이관해 별도로 보정 처리한다.

Scrapy: Strategy + Chain of Responsibility

# spiders/base.py
from typing import Protocol

class ParseStrategy(Protocol):
    def parse(self, response) -> dict: ...

# spiders/strategies.py
class SiteAStrategy:
    def parse(self, response):
        return {
            "sku": response.css("#sku::text").get(),
            "name": response.css("h1::text").get(),
            "price": response.css(".price::text").re_first(r"\d+(\.\d+)?"),
        }

class SiteBStrategy:
    def parse(self, response):
        # 사이트 B 전용 규칙
        return {
            "sku": response.xpath("//span[@id='sku']/text()").get(),
            "name": response.xpath("//h1/text()").get(),
            "price": response.xpath("//span[@class='p']/text()").re_first(r"\d+(\.\d+)?"),
        }

# spiders/product_spider.py
import scrapy
from .strategies import SiteAStrategy

class ProductSpider(scrapy.Spider):
    name = "products_a"
    start_urls = ["https://example.com/products/1"]
    strategy = SiteAStrategy()

    def parse(self, response):
        yield self.strategy.parse(response)
# pipelines.py
from itemadapter import ItemAdapter
from tasks import upsert_product_task

class ValidatePipeline:
    def process_item(self, item, spider):
        a = ItemAdapter(item)
        assert a.get("sku") and a.get("name") and a.get("price"), "필수 필드 누락"
        return item

class NormalizePipeline:
    def process_item(self, item, spider):
        a = ItemAdapter(item)
        a["name"] = a["name"].strip()
        a["price"] = float(str(a["price"]).replace(",", ""))
        return item

class DispatchPipeline:
    def process_item(self, item, spider):
        upsert_product_task.delay(dict(item))  # 비동기 커맨드
        return item
# settings.py (Scrapy)
ITEM_PIPELINES = {
    "pipelines.ValidatePipeline": 100,
    "pipelines.NormalizePipeline": 200,
    "pipelines.DispatchPipeline": 300,
}

파이프라인에서는 장기 실행 로직(네트워크·DB 직접 접근)을 피하고 큐로 위임하는 것이 원칙이다. 사이트별 Strategy는 Spider 교체 없이 조립 방식으로 확장한다.

계층별로 본 패턴의 특성

패턴 적용 계층 성능 확장성 일관성 운영 편의
Repository + UoW SQLModel/Service 트랜잭션 최소화로 안정적, 약간의 오버헤드 ORM 교체·테스트 용이 강한 일관성 확보 마이그레이션/테스트 편리
Dependency Injection FastAPI 호출당 객체 생성 비용 미미 조립 변경 용이 구성 일관성 설정/테스트 단순
Strategy Scrapy Spider 런타임 분기 최소 신규 사이트 추가 용이 파싱 규칙 분리 코드 가독성 향상
Chain of Responsibility Scrapy Pipelines 단계별 단순 처리 단계 추가/변경 용이 검증·정규화 표준화 모니터링 포인트 명확
Command + Retry + Idempotency Celery Tasks 재시도 비용 존재 워커 스케일아웃 용이 멱등키로 결과 안정 장애 복구 자동화

크롤링부터 배치 적재까지

상품 크롤링 → 정규화 → 비동기 upsert → API 제공 흐름에서는 Scrapy가 사이트별 Strategy로 파싱하고 Pipeline이 큐에 태스크를 전송한다. Celery Worker가 UoW로 upsert하고 FastAPI가 조회 API를 제공하며, 재시도·멱등성으로 운영 안정성을 확보한다.

가격 변경 이벤트 발행에서는 UoW 커밋 직후 Outbox 레코드를 생성하고 별도 Publisher가 전송한다(Transactional Outbox). 재전송에 대비해 상태를 관리하고, 외부 소비자(검색·캐시)에게 전달해 읽기 성능을 높인다.

대량 적재 배치에서는 크롤링 결과를 청크 단위로 태스크 분할(Command)하고 워커 풀을 확장한다. Repository에 벌크 upsert 전략을 추가하고, 고유 제약을 이용해 충돌을 최소화한다. 장애 시에는 재처리가 자동화되고 DLQ 분리로 수동 개입이 최소화되며, 구성 가능한 DI·UoW로 로컬·스테이징·운영 환경 전환도 쉬워진다.

운영 시 주의할 점

긴 트랜잭션은 락 경합을 유발할 수 있어 UoW 범위를 최소화하는 것이 좋다. 무분별한 재시도는 쓰기 폭증으로 이어질 수 있으므로 백오프·지수·Jitter를 적용하고 한도를 설정해야 한다. Outbox·이벤트 발행은 최종적 일관성이라는 트레이드오프가 있어 소비자 측의 재처리·멱등 설계가 필요하다.

초기에는 UoW, DI, Chain 같은 최소 패턴부터 적용하고, 크기와 복잡도가 늘면 Strategy·Outbox로 확장하는 점진적 도입이 권장된다.

Repository 패턴Unit of WorkStrategy 패턴Command 패턴FastAPI