diff --git a/tools/import_thirdparty_catalog_csv.py b/tools/import_thirdparty_catalog_csv.py index 4ac1212..5385e1d 100644 --- a/tools/import_thirdparty_catalog_csv.py +++ b/tools/import_thirdparty_catalog_csv.py @@ -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()