refactor(negodata): LPS 자동 요청 배치 제거 — 검색은 수동 트리거 전용
협의 결정(2026-07-10): 크롤 비용이 사용자 행동에만 비례하도록,
갱신 오래된 상품을 자동으로 검색 요청하던 잡(매일 04:00, 상한 500건,
일 ~$2)을 제거한다. 검색 진입점은 상품 화면의 수동 트리거
(POST /v1/item/{id}/lowest-price) 하나만 남는다.
- request_lps_searches 잡·request_stale_searches 서비스·stale_items
CRUD·배치 상수(REFRESH_HOURS 등) 제거
- 수집 잡(sync_lps_results, 5분)은 유지 — 크롤을 일으키지 않는
반영 백스톱(모달 조기 종료·폴링 초과분 자동 반영, 비용 0)
검증: 스케줄러 등록 잡 3개(견적마감 2 + LPS 수집 1) 확인, 수집 잡 실행 정상
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
parent
3171311746
commit
5ddf5f6249
@ -10,7 +10,7 @@
|
||||
from abc import ABC, abstractmethod
|
||||
from typing import Tuple
|
||||
|
||||
from sqlalchemy import column, func, or_, select, table, update
|
||||
from sqlalchemy import column, func, select, table, update
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from common.database.db_session_manager import DB_SESSION_MNG
|
||||
@ -34,10 +34,6 @@ class ILpsSyncCRUD(ABC):
|
||||
async def watermark(self, cdb: AsyncSession) -> Tuple[ErrorType, object]:
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
async def stale_items(self, cdb: AsyncSession, cutoff, limit) -> Tuple[ErrorType, list]:
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
async def existing_item_ids(self, cdb: AsyncSession, item_ids: list) -> Tuple[ErrorType, set]:
|
||||
pass
|
||||
@ -72,34 +68,6 @@ class LpsSyncCRUD(ILpsSyncCRUD):
|
||||
LOG.e_no_callstack(f"[lps-sync] watermark 조회 실패: {ex}")
|
||||
return ErrorType.DB_RUN_FAILED, None
|
||||
|
||||
async def stale_items(self, cdb: AsyncSession, cutoff, limit: int):
|
||||
"""재검색 대상 상품 — 수집 이력이 없거나 마지막 수집이 cutoff 이전인 활성 상품.
|
||||
반환 행: (item_id, name, model_name, spec, manufacturer, price)"""
|
||||
try:
|
||||
latest = (
|
||||
select(
|
||||
item_internet_lowest_prices.item_id.label("item_id"),
|
||||
func.max(item_internet_lowest_prices.crawl_end_time).label("last_ts"),
|
||||
)
|
||||
.where(item_internet_lowest_prices.deleted == False) # noqa: E712
|
||||
.group_by(item_internet_lowest_prices.item_id)
|
||||
.subquery()
|
||||
)
|
||||
q = (
|
||||
select(items.item_id, items.name, items.model_name, items.spec, items.manufacturer, items.price)
|
||||
.join(latest, items.item_id == latest.c.item_id, isouter=True)
|
||||
.where(
|
||||
items.deleted == False, # noqa: E712
|
||||
or_(latest.c.last_ts.is_(None), latest.c.last_ts < cutoff),
|
||||
)
|
||||
.order_by(latest.c.last_ts.asc().nullsfirst()) # 오래된 것부터(이력 없는 신규 최우선)
|
||||
.limit(limit)
|
||||
)
|
||||
return ErrorType.SUCCESS, (await cdb.execute(q)).all()
|
||||
except Exception as ex:
|
||||
LOG.e_no_callstack(f"[lps-sync] stale_items 조회 실패: {ex}")
|
||||
return ErrorType.DB_RUN_FAILED, []
|
||||
|
||||
async def existing_item_ids(self, cdb: AsyncSession, item_ids: list):
|
||||
"""전달된 id 중 실재하는 활성 상품 id 집합 — 결과 반영 전 매핑 검증용."""
|
||||
if not item_ids:
|
||||
|
||||
@ -4,8 +4,11 @@
|
||||
|
||||
잡 ① close_expired_quotations : 5분마다(KST) — 마감시각 지난 견적 마감
|
||||
잡 ② close_negotiated_quotations: 5분마다(KST) — 모든 세션 협상 끝난 견적 즉시 마감(타입 무관)
|
||||
잡 ③ request_lps_searches : 매일 04:00(KST) — 갱신 오래된 상품을 LPS 에 검색 요청(비용 발생 잡)
|
||||
잡 ④ sync_lps_results : 5분마다(KST) — lps_db 결과 증분 수집 → 이력 append + 상품 최저가 박제
|
||||
잡 ③ sync_lps_results : 5분마다(KST) — lps_db 결과 증분 수집 → 이력 append + 상품 최저가 박제
|
||||
(수동 트리거된 검색의 반영 백스톱 — 크롤을 일으키지 않음, 비용 0)
|
||||
|
||||
LPS 검색 **요청**은 배치로 돌리지 않는다(2026-07-10 협의) — 상품 화면의 수동 트리거
|
||||
(POST /v1/item/{id}/lowest-price)로만 검색한다. 크롤 비용이 사용자 행동에만 비례하게.
|
||||
"""
|
||||
import os
|
||||
|
||||
@ -53,16 +56,7 @@ def start_scheduler():
|
||||
misfire_grace_time=600,
|
||||
max_instances=1,
|
||||
)
|
||||
# 잡 ③ LPS 검색 요청(비용 발생) — 새벽 1회. LPS 미설정 환경이면 잡 내부에서 스킵.
|
||||
_scheduler.add_job(
|
||||
jobs.request_lps_searches,
|
||||
CronTrigger(hour=4, minute=0),
|
||||
id="request_lps_searches",
|
||||
coalesce=True,
|
||||
misfire_grace_time=3600, # 재기동 등으로 놓쳐도 1시간 내면 실행
|
||||
max_instances=1,
|
||||
)
|
||||
# 잡 ④ LPS 결과 수집·반영 — 증분·멱등이라 잦아도 안전
|
||||
# 잡 ③ LPS 결과 수집·반영 — 증분·멱등이라 잦아도 안전. 크롤 요청은 하지 않는다(수동 트리거 전용).
|
||||
_scheduler.add_job(
|
||||
jobs.sync_lps_results,
|
||||
CronTrigger(minute="*/5"),
|
||||
@ -72,7 +66,7 @@ def start_scheduler():
|
||||
max_instances=1,
|
||||
)
|
||||
_scheduler.start()
|
||||
LOG.i("[scheduler] started (KST: 견적마감 2잡 5분 · LPS 요청 04:00 · LPS 수집 5분)")
|
||||
LOG.i("[scheduler] started (KST: 견적마감 2잡 5분 · LPS 수집 5분 — LPS 요청 배치 없음, 수동 트리거 전용)")
|
||||
|
||||
|
||||
def shutdown_scheduler():
|
||||
|
||||
@ -83,26 +83,9 @@ async def close_negotiated_quotations() -> int:
|
||||
|
||||
# ---- LPS(인터넷 최저가) 동기화 ------------------------------------------
|
||||
|
||||
async def request_lps_searches() -> int:
|
||||
"""[잡③] 갱신이 오래된 상품을 LPS 에 검색 요청(enqueue). 매일 새벽 1회.
|
||||
비용이 발생하는 잡(상품당 ~$0.004) — 주기·상한은 lps_sync_service 상수로 관리.
|
||||
LPS 미설정 환경이면 조용히 스킵(available=False)."""
|
||||
from services.lps_sync_service import LpsSyncService
|
||||
|
||||
service = LpsSyncService()
|
||||
if not service.available():
|
||||
return 0
|
||||
results = await service.request_stale_searches()
|
||||
if results:
|
||||
LOG.i(
|
||||
f"[scheduler] lps_request: 접수 {results['accepted']} / 활성중복 {results['duplicated']} / "
|
||||
f"이름없음 {results['skipped_no_name']} / HTTP오류 {results['http_error']}"
|
||||
)
|
||||
return results["accepted"]
|
||||
|
||||
|
||||
async def sync_lps_results() -> int:
|
||||
"""[잡④] lps_db.price_history 증분을 읽어 수집 이력 append + 상품 대표 최저가 박제. 5분마다.
|
||||
"""[잡③] lps_db.price_history 증분을 읽어 수집 이력 append + 상품 대표 최저가 박제. 5분마다.
|
||||
검색 **요청**은 하지 않는다(수동 트리거 전용, 2026-07-10 협의) — 이 잡은 반영 백스톱.
|
||||
워터마크(=이력의 max crawl_end_time) 기준 증분이라 멱등 — 실패 tick 은 다음 tick 이 흡수."""
|
||||
from services.lps_sync_service import LpsSyncService
|
||||
|
||||
|
||||
@ -1,10 +1,13 @@
|
||||
"""LPS(인터넷 최저가 검색) 동기화 서비스 — 요청·수집·반영의 도메인 로직.
|
||||
|
||||
흐름 (요청은 API, 결과는 DB — 2026-07-10 설계 결정)
|
||||
① 요청: 갱신이 오래된 상품을 골라 LPS API(POST /v1/lps/search)로 enqueue.
|
||||
product_code = items.item_id(uuid 문자열) — 이걸로 결과가 자동 매핑된다.
|
||||
enqueue 를 DB insert 로 하지 않는 이유: LPS 의 활성중복 dedupe·pg_notify 워커 깨움을 우회하게 됨.
|
||||
① 요청: **수동 트리거 전용**(상품 화면 → POST /v1/item/{id}/lowest-price → 여기의
|
||||
request_search_for_item). 자동 요청 배치는 두지 않는다(2026-07-10 협의 — 크롤 비용이
|
||||
사용자 행동에만 비례하게). product_code = items.item_id(uuid 문자열) — 이걸로 결과가
|
||||
자동 매핑된다. enqueue 를 DB insert 로 하지 않는 이유: LPS 의 활성중복 dedupe·pg_notify
|
||||
워커 깨움을 우회하게 됨.
|
||||
② 수집: lps_db.price_history(읽기전용)를 워터마크(max crawl_end_time) 증분으로 읽는다.
|
||||
(5분 크론 + 조회 API 의 온디맨드 — 트리거된 검색의 반영 백스톱, 크롤 비용 0)
|
||||
③ 반영: 이력은 partner.item_internet_lowest_prices 에 append(성공/실패 모두),
|
||||
성공분의 상품별 최신값을 items.internet_lowest_price 에 박제(+yn).
|
||||
|
||||
@ -27,11 +30,6 @@ from common.logger import LOG
|
||||
from config.server_configs import web_server_config
|
||||
from crud.lps_sync_crud import ILpsSyncCRUD, LpsSyncCRUD
|
||||
|
||||
# 재검색 주기·요청 배치 크기 기본값. 상품당 실측 비용 ~$0.004 이므로
|
||||
# (상품 수 × 24h 주기)로 월 비용이 바로 계산된다. 필요 시 여기만 조정.
|
||||
REFRESH_HOURS = 24 # 마지막 수집이 이보다 오래된 상품만 재요청
|
||||
REQUEST_BATCH_LIMIT = 500 # 요청 잡 1회가 enqueue 하는 최대 상품 수(비용 상한)
|
||||
ENQUEUE_CHUNK = 100 # LPS API 1콜에 담는 상품 수
|
||||
SCAN_FLOOR_DAYS = 3 # price_history 스캔 하한(일) — 스킵행(비uuid 등)은 워터마크를
|
||||
# 못 올리므로, 바닥 없이는 같은 행을 영원히 재스캔하게 된다
|
||||
|
||||
@ -58,52 +56,6 @@ class LpsSyncService:
|
||||
"""LPS 연동 활성 여부 — lps_db 가 등록된 환경에서만 배치가 돈다."""
|
||||
return DB_SESSION_MNG.is_registered(DBType.LPS.value)
|
||||
|
||||
# ---- ① 요청: 오래된 상품을 LPS 에 검색 enqueue ----------------------
|
||||
async def request_stale_searches(self) -> Counter:
|
||||
results = Counter()
|
||||
if not self.available():
|
||||
return results
|
||||
|
||||
cutoff = _utc_now_aware() - timedelta(hours=REFRESH_HOURS) # timestamptz 비교 — aware 필수(위 sync_results 주석)
|
||||
|
||||
err, rows = await DB_SESSION_MNG.execute_lambda(
|
||||
DBType.MAIN.value, DBWRType.DB_READ.value,
|
||||
lambda s: self.crud.stale_items(s, cutoff, REQUEST_BATCH_LIMIT),
|
||||
)
|
||||
if err != ErrorType.SUCCESS or not rows:
|
||||
return results
|
||||
|
||||
payload_items = []
|
||||
for item_id, name, model_name, spec, manufacturer, price in rows:
|
||||
if not (name or "").strip():
|
||||
results["skipped_no_name"] += 1
|
||||
continue
|
||||
payload_items.append({
|
||||
"product_code": str(item_id),
|
||||
"product_name": name,
|
||||
"job_type": "batch",
|
||||
"model": model_name or "",
|
||||
"specification": spec or "",
|
||||
"company": manufacturer or "",
|
||||
"price": str(price) if price else "",
|
||||
})
|
||||
|
||||
base = web_server_config.lps_base_url.rstrip("/")
|
||||
async with httpx.AsyncClient(timeout=10.0) as client:
|
||||
for i in range(0, len(payload_items), ENQUEUE_CHUNK):
|
||||
chunk = payload_items[i:i + ENQUEUE_CHUNK]
|
||||
try:
|
||||
r = await client.post(f"{base}/v1/lps/search", json={"data": chunk})
|
||||
r.raise_for_status()
|
||||
body = r.json()
|
||||
results["accepted"] += int(body.get("accepted", 0))
|
||||
results["duplicated"] += sum(1 for it in body.get("items", []) if it.get("duplicated"))
|
||||
except Exception as ex:
|
||||
# LPS 다운은 이 배치의 실패일 뿐 — 다음 tick 이 같은 대상(여전히 stale)을 다시 요청한다.
|
||||
results["http_error"] += 1
|
||||
LOG.w(f"[lps-sync] 검색요청 실패(청크 {i // ENQUEUE_CHUNK}): {type(ex).__name__}: {ex}")
|
||||
return results
|
||||
|
||||
# ---- 단건 즉시 요청 (lowest-price 트리거 API 용) ---------------------
|
||||
async def request_search_for_item(self, item) -> tuple:
|
||||
"""상품 1건을 즉시 LPS 에 검색 요청(수동 트리거 — job_type=manual, 배치보다 높은 우선순위).
|
||||
|
||||
Loading…
Reference in New Issue
Block a user