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% 이상 감소할 것으로 기대된다.
데이터가 흐르는 경로
입력은 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로 확장하는 점진적 도입이 권장된다.