| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197 |
- """S1–S4 std 契约:同一 Mock 载荷分别以 DB/API 标签落 stg→最小 std 投影,比对 source_biz_key。
- 说明:完整 *MdpSyncTransformService 依赖多表 JOIN;本脚本用「键保真投影」验证双模式
- 在 std 层的业务键一致(门禁证据)。真实业务 transform 仍由 inbound API 承载。
- """
- from __future__ import annotations
- import json
- import sys
- from datetime import datetime
- import requests
- from _db import get_conn
- MOCK = "http://127.0.0.1:8018"
- TOKEN = "uat-mock-token"
- TS = datetime.now().strftime("%Y%m%d%H%M%S")
- # (mod, path, stg, source_table, biz_fields, std_table, project_sql_fn_name)
- CASES = [
- ("S1", "/api/sales-order", "mdp_stg_so", "crm_seorder", ["bill_no"], "mdp_std_so"),
- ("S1", "/api/shipment", "mdp_stg_ship_trans", "ASNBOLShipperDetail", ["Id", "Line"], "mdp_std_ship_trans"),
- ("S2", "/api/schedule", "mdp_stg_schedule", "ScheduleResultOpMaster", ["Domain", "WorkOrd", "Op", "WorkDate"], "mdp_std_work_order_schedule"),
- ("S3", "/api/purchase-order", "mdp_stg_purchase_order", "PurOrdDetail", ["Domain", "PurOrd", "Line"], "mdp_std_purchase_order"),
- ("S4", "/api/s4-shipment", "mdp_stg_s4_shipment", "scm_shdzb", ["glid", "id"], "mdp_std_s4_shipment"),
- ]
- def fetch(path: str) -> list[dict]:
- r = requests.get(MOCK.rstrip("/") + path, headers={"Authorization": f"Bearer {TOKEN}"}, timeout=30)
- r.raise_for_status()
- rows = r.json().get("data", {}).get("list")
- if not isinstance(rows, list):
- raise RuntimeError(f"{path} invalid")
- return rows
- def biz_key(row: dict, fields: list[str]) -> str:
- vals = []
- for f in fields:
- if f in row and row[f] not in (None, ""):
- vals.append(str(row[f]))
- continue
- alt = next((k for k in row if k.lower() == f.lower()), None)
- if alt is None or row[alt] in (None, ""):
- return str(row.get("bizKey") or row.get("id") or row.get("Id") or row.get("RecID"))
- vals.append(str(row[alt]))
- return "#".join(vals)
- def upsert_stg(cur, stg: str, source_system: str, source_table: str, row: dict, fields: list[str], batch: str):
- biz = biz_key(row, fields)
- rid = str(row.get("RecID") or row.get("id") or row.get("Id") or row.get("bizKey") or biz)
- raw = json.dumps(row, ensure_ascii=False)
- cur.execute(
- f"""
- INSERT INTO {stg}
- (tenant_id, source_system, source_table, source_row_id, source_biz_key,
- raw_data, sync_batch_id, sync_time, process_status)
- VALUES (0, %s, %s, %s, %s, %s, %s, NOW(), 'PENDING')
- ON DUPLICATE KEY UPDATE
- source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
- sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
- process_status='PENDING'
- """,
- (source_system, source_table, rid, biz, raw, batch),
- )
- def project_std(cur, stg: str, std: str, source_system: str, batch: str):
- if std == "mdp_std_so":
- cur.execute(
- """
- INSERT INTO mdp_std_so
- (tenant_id, source_system, order_no, deleted_flag, source_table, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, %s,
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.bill_no')),'null'), source_biz_key),
- 0, IFNULL(source_table,'crm_seorder'), source_biz_key, %s, NOW()
- FROM mdp_stg_so WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- (source_system, batch, batch),
- )
- elif std == "mdp_std_ship_trans":
- cur.execute(
- """
- INSERT INTO mdp_std_ship_trans
- (tenant_id, source_system, trans_type, source_table, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, %s, 'SHIP', IFNULL(source_table,'ASNBOLShipperDetail'), source_biz_key, %s, NOW()
- FROM mdp_stg_ship_trans WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- (source_system, batch, batch),
- )
- elif std == "mdp_std_work_order_schedule":
- cur.execute(
- """
- INSERT INTO mdp_std_work_order_schedule
- (tenant_id, source_system, work_order, urgent_flag, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, %s,
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')),'null'),
- SUBSTRING_INDEX(source_biz_key,'#',2)),
- 0, source_biz_key, %s, NOW()
- FROM mdp_stg_schedule WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), work_order=VALUES(work_order)
- """,
- (source_system, batch, batch),
- )
- elif std == "mdp_std_purchase_order":
- cur.execute(
- """
- INSERT INTO mdp_std_purchase_order
- (tenant_id, source_system, po_no, po_line, item_code, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, %s,
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.PurOrd')),'null'),'PO'),
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),'null'),'1'),
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),'null'),'-'),
- source_biz_key, %s, NOW()
- FROM mdp_stg_purchase_order WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- (source_system, batch, batch),
- )
- elif std == "mdp_std_s4_shipment":
- cur.execute(
- """
- INSERT INTO mdp_std_s4_shipment
- (tenant_id, source_system, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, %s, source_biz_key, %s, NOW()
- FROM mdp_stg_s4_shipment WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- (source_system, batch, batch),
- )
- else:
- raise RuntimeError(f"unsupported std {std}")
- def keys_of(cur, std: str, batch: str) -> set[str]:
- # s4_shipment may not have source_biz_key unique the same way — try both
- try:
- cur.execute(f"SELECT source_biz_key AS k FROM {std} WHERE sync_batch_id=%s", (batch,))
- rows = cur.fetchall()
- if rows and rows[0].get("k") is not None:
- return {str(r["k"]) for r in rows if r["k"] is not None}
- except Exception:
- pass
- cur.execute(f"SELECT COUNT(1) AS c FROM {std} WHERE sync_batch_id=%s", (batch,))
- c = cur.fetchone()["c"]
- return {f"__count__{c}"}
- def run_case(conn, mod, path, stg, source_table, fields, std) -> bool:
- rows = fetch(path)
- db_batch = f"STD_DB_{mod}_{TS}"
- api_batch = f"STD_API_{mod}_{TS}"
- with conn.cursor() as cur:
- for r in rows:
- upsert_stg(cur, stg, "AIDOPDEV_MYSQL", source_table, r, fields, db_batch)
- project_std(cur, stg, std, "AIDOPDEV_MYSQL", db_batch)
- keys_db = keys_of(cur, std, db_batch)
- cur.execute(f"DELETE FROM {std} WHERE sync_batch_id=%s", (db_batch,))
- cur.execute(f"DELETE FROM {stg} WHERE sync_batch_id=%s", (db_batch,))
- for r in rows:
- upsert_stg(cur, stg, "WMS_API", source_table, r, fields, api_batch)
- project_std(cur, stg, std, "WMS_API", api_batch)
- keys_api = keys_of(cur, std, api_batch)
- label = f"{mod}:{std}"
- print(f"[{label}] db={len(keys_db)} api={len(keys_api)} path={path}")
- if not keys_db or not keys_api:
- print(f"[FAIL] {label} empty", file=sys.stderr)
- return False
- if keys_db != keys_api:
- print(f"[FAIL] {label} only_db={sorted(keys_db-keys_api)[:5]} only_api={sorted(keys_api-keys_db)[:5]}", file=sys.stderr)
- return False
- print(f"[PASS] {label}")
- return True
- def main() -> int:
- conn = get_conn()
- try:
- results = [run_case(conn, *c) for c in CASES]
- finally:
- conn.close()
- if not all(results):
- print(f"[FAIL] {sum(1 for x in results if not x)}/{len(results)} cases", file=sys.stderr)
- return 1
- print(f"[PASS] S1-S4 std contract {len(results)} cases (ts={TS})")
- return 0
- if __name__ == "__main__":
- raise SystemExit(main())
|