Files
shaguabijia-app-server/app/repositories/store_mapping.py
T
marco fdd3d8cf9b feat(store_mapping): 加京东原料列 + upsert 首写泛化(任意平台首写不丢列)
京东秒送接入: 店铺身份是两个稳定数字 id —— storeId 进 id_jd(同 taobao shopId→id_taobao),
venderId(deeplink 还需、商家维度)单列存 jd_vender_id; 另加 jd_share_url/jd_resolved_url/
jd_deeplink 原料列(与 taobao_*/meituan_* 平行)。迁移 store_mapping_jd_cols 加这 4 列 +
jd_vender_id 索引。

顺手修一个潜在 bug: upsert 的首写(INSERT)原来硬编码只灌 taobao+通用列, 首写的若是美团/
京东腿, 它们的列不在构造器里会被丢。改成按 _MERGE_COLUMNS 逐列 setattr, 任意平台首写都全量
落库(美团也受益)。合并(填空不覆盖)逻辑不变。

验证: 迁移单 head、store_mapping 无新漂移; TestClient 打真 PG —— 京东首写 4 列全落 →
淘宝同 trace 合并填淘宝列、不覆盖京东列。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-13 22:35:39 +08:00

91 lines
4.0 KiB
Python

"""平台店铺映射落库:一次比价(trace)一行,各目标平台解析出 id 就 upsert 进同一行。
合并键 = trace_id(唯一约束)。一次跨平台比价 = 一个 trace = 一个真实店铺:淘宝腿解析出
shopId 先建行(填淘宝列),京东腿(将来)解析出 id 再 upsert 进**同一行**(填京东列)。
合并策略 = 填空(fill-the-blanks):只写该行当前为 NULL 的列,绝不覆盖已有非空值。保证后到
的平台只填自己那几列、动不了先到平台的数据;共享列(geo / source / 溯源)先到先得。
不依赖 ON CONFLICT,跨方言(PG / SQLite dev)都安全;并发撞唯一约束则回滚后转走合并路径。
⚠️ 一行里 id_taobao 与 id_jd 共存只表示"两条腿搜同一个源店名各自匹配到了某家店",是
name-match 置信度、非已核实同一实体。作为 append-only 原始资产留存,跨平台精确匹配由下游做。
"""
from __future__ import annotations
import logging
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.models.store_mapping import StoreMapping
from app.schemas.store_mapping import StoreMappingIn
logger = logging.getLogger("shagua.store_mapping")
# upsert 填空时可写的列。不含: trace_id(合并键)/ business_type(非空默认)/
# id(主键)/ created_at(server_default)。各平台只会带自己那几列非空, 其余为 None 不动。
_MERGE_COLUMNS = (
"source_platform",
"id_taobao", "name_taobao", "id_meituan", "name_meituan", "id_jd", "name_jd",
"city", "geohash", "lng", "lat", "taobao_address",
"source_device_id", "source_user_id",
"taobao_share_url", "taobao_resolved_url", "taobao_deeplink",
"meituan_poi_id_str", "meituan_share_url", "meituan_resolved_url", "meituan_deeplink",
"jd_vender_id", "jd_share_url", "jd_resolved_url", "jd_deeplink",
"attrs",
)
def _find(db: Session, trace_id: str) -> StoreMapping | None:
return db.execute(
select(StoreMapping).where(StoreMapping.trace_id == trace_id)
).scalar_one_or_none()
def _merge_fill_blanks(existing: StoreMapping, payload: StoreMappingIn) -> list[str]:
"""把 payload 里非空、且 existing 当前为 NULL 的列填进去。返回被填的列名(空=无变化)。"""
filled: list[str] = []
for col in _MERGE_COLUMNS:
new = getattr(payload, col)
if new is not None and getattr(existing, col) is None:
setattr(existing, col, new)
filled.append(col)
return filled
def upsert(db: Session, payload: StoreMappingIn) -> tuple[int, int | None]:
"""落一行跨平台店铺映射,返回 (created, row_id)。
created=1 新建该 trace 的行 / 0 合并进已存在行(填空,不覆盖)。"""
existing = _find(db, payload.trace_id)
if existing is None:
# 首写: 把 payload 全部可合并列灌进去(逐列 setattr 而非硬编码构造器, 否则首写的是
# 美团/京东腿时它们的列会漏 —— 不在硬编码列表里就丢)。trace_id/business_type 是键/默认, 显式给。
row = StoreMapping(
trace_id=payload.trace_id,
business_type=payload.business_type,
)
for col in _MERGE_COLUMNS:
setattr(row, col, getattr(payload, col))
db.add(row)
try:
db.commit()
except IntegrityError:
# 并发: 另一个请求刚插了同 trace → 撞唯一约束。回滚后转合并路径填空。
db.rollback()
existing = _find(db, payload.trace_id)
if existing is None:
raise
logger.warning("store_mapping 并发冲突 trace=%s, 转填空合并", payload.trace_id)
else:
db.refresh(row)
return 1, row.id
# 已存在(或并发回退到此): 填空合并, 只写当前 NULL 的列
filled = _merge_fill_blanks(existing, payload)
if filled:
db.commit()
db.refresh(existing)
logger.info("store_mapping 合并 trace=%s 填列=%s", payload.trace_id, filled)
return 0, existing.id