fix: 跳过已确认保留的目录关联 (#281)

This commit is contained in:
chengma
2026-08-20 17:15:18 +08:00
parent efad2eadf4
commit ac2250a959
+21 -8
View File
@@ -32,6 +32,7 @@ def product_from(row):
def sku_from(row):
return {'sku_id':row['shopee_sku_id'],'goods_id':row['shopee_goods_id'],'spec_raw':row['spec_raw'],'color':row['color'],'size':row['size'],'advice':row['advice'],'parse_ok':row['parse_ok']=='true','sku_code':row['shopee_sku_code']}
def pdd_from(row): return {'goods_id':row['pdd_goods_id'],'url':row['pdd_goods_url'],'title':'','shop_name':'','dimensions':[],'skus':[]}
def spec_key(raw): return ' '.join(raw.split())
def groups(root):
"""流式读 CSV;来源按商品连续导出,同商品不会一次占用整个文件。"""
@@ -49,19 +50,30 @@ def groups(root):
def make_group(rows):
first=rows[0]; skus=[]; specs=set(); pdds={}; associations=[]
for row in rows:
if row['spec_raw'] and row['spec_raw'] not in specs:
specs.add(row['spec_raw']); skus.append(sku_from(row))
key=spec_key(row['spec_raw'])
if key and key not in specs:
specs.add(key); skus.append(sku_from(row))
if row['pdd_importable']=='true' and row['pdd_goods_id'] not in pdds:
pdds[row['pdd_goods_id']]=pdd_from(row)
associations.append({'shopee_goods_id':first['shopee_goods_id'],'pdd_goods_id':row['pdd_goods_id']})
return product_from(first),skus,pdds,associations
def batches(root, observed_at, dry_run):
def batches(root, observed_at, dry_run, keep_existing):
selected=[]; seen_pdd={}; sku_count=0
def emit():
nonlocal selected,seen_pdd,sku_count
products=[g[0] for g in selected]; skus=[s for g in selected for s in g[1]]
pdds=list(seen_pdd.values()); assocs=[a for g in selected for a in g[3]]
# 同一商品可能在不同源 CSV 再次出现;接口禁止批内重复身份,
# 所以在真正组装请求时按四类业务键再去重一次。
product_by_id={g[0]['goods_id']:g[0] for g in selected}
sku_by_key={}
association_by_key={}
for group in selected:
for sku in group[1]: sku_by_key.setdefault((sku['goods_id'],spec_key(sku['spec_raw'])),sku)
for association in group[3]:
if association['shopee_goods_id'] not in keep_existing:
association_by_key.setdefault((association['shopee_goods_id'],association['pdd_goods_id']),association)
products=list(product_by_id.values()); skus=list(sku_by_key.values())
pdds=list(seen_pdd.values()); assocs=list(association_by_key.values())
seed=json.dumps([products,skus,pdds,assocs],ensure_ascii=False,sort_keys=True,separators=(',',':')).encode()
payload={'schema_version':1,'batch_id':'thirdparty-csv-v1-'+hashlib.sha256(seed).hexdigest()[:24],'observed_at':observed_at,'update_policy':'fill_missing','dry_run':dry_run,'shopee_products':products,'shopee_skus':skus,'pdd_products':pdds,'associations':assocs}
if len(json.dumps(payload,ensure_ascii=False,separators=(',',':')).encode()) > MAX_BYTES:
@@ -75,18 +87,19 @@ def batches(root, observed_at, dry_run):
if selected: yield emit()
def main():
p=argparse.ArgumentParser(); p.add_argument('csv_dir'); p.add_argument('--base-url',required=True); p.add_argument('--observed-at',default='2026-08-20T00:00:00+08:00'); mode=p.add_mutually_exclusive_group(required=True); mode.add_argument('--dry-run',action='store_true'); mode.add_argument('--apply',action='store_true'); p.add_argument('--report',required=True); a=p.parse_args()
p=argparse.ArgumentParser(); p.add_argument('csv_dir'); p.add_argument('--base-url',required=True); p.add_argument('--observed-at',default='2026-08-20T00:00:00+08:00'); mode=p.add_mutually_exclusive_group(required=True); mode.add_argument('--dry-run',action='store_true'); mode.add_argument('--apply',action='store_true'); p.add_argument('--report',required=True); p.add_argument('--keep-existing-association',action='append',default=[],help='保留该蝦皮商品的数据库既有关联,不写来源新关联;可重复传入'); a=p.parse_args()
token=os.environ.get('CMAUTOBUY_CATALOG_TOKEN','').strip()
if not token: error('缺少环境变量 CMAUTOBUY_CATALOG_TOKEN')
root=Path(a.csv_dir)
endpoint=a.base_url.rstrip('/')+'/api/v1/integrations/catalog/batches'
total=Counter(); conflicts=[]; count=0
for payload in batches(root,a.observed_at,a.dry_run):
keep_existing={value.strip() for value in a.keep_existing_association if value.strip()}
for payload in batches(root,a.observed_at,a.dry_run,keep_existing):
count+=1; response=request(endpoint,token,payload)
if response.get('status') not in {'previewed','succeeded'}: error('接口未返回预期状态')
total.update(response.get('counts',{})); conflicts.extend(response.get('conflicts',[]))
print(f"[{count}] {response['status']} {payload['batch_id']}")
Path(a.report).write_text(json.dumps({'mode':'preview' if a.dry_run else 'apply','batches':count,'counts':total,'association_conflicts':conflicts},ensure_ascii=False,indent=2),encoding='utf-8')
Path(a.report).write_text(json.dumps({'mode':'preview' if a.dry_run else 'apply','batches':count,'counts':total,'kept_existing_association_goods_ids':sorted(keep_existing),'association_conflicts':conflicts},ensure_ascii=False,indent=2),encoding='utf-8')
print(f"完成:{count} 批;报告:{a.report}")
if __name__=='__main__':
try: main()