run_contract_s5_s6_s7.py 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309
  1. """S5/S6/S7 契约:同一 Mock 载荷分别以 DB/API 源标签落 stg→std,比对 source_biz_key 集合。
  2. 说明:mdp_std_* 业务唯一键在 DB/API 两路上会冲突,后写会覆盖 sync_batch_id,
  3. 故不能依赖「两批同时留在 std 再 LIKE 比对」。本脚本改为:
  4. 1) 跑 DB 路线 → 快照键集 A → 删除该批 std
  5. 2) 跑 API 路线 → 快照键集 B
  6. 3) 断言 A≡B(行数与键集合一致)
  7. 用法:
  8. set MOCK_TOKEN=uat-mock-token
  9. python run_contract_s5_s6_s7.py
  10. """
  11. from __future__ import annotations
  12. import json
  13. import sys
  14. from datetime import datetime
  15. import requests
  16. from _db import get_conn
  17. MOCK = "http://127.0.0.1:8018"
  18. TOKEN = "uat-mock-token"
  19. TS = datetime.now().strftime("%Y%m%d%H%M%S")
  20. def fetch(path: str) -> list[dict]:
  21. r = requests.get(
  22. MOCK.rstrip("/") + path,
  23. headers={"Authorization": f"Bearer {TOKEN}"},
  24. timeout=30,
  25. )
  26. r.raise_for_status()
  27. rows = r.json().get("data", {}).get("list")
  28. if not isinstance(rows, list):
  29. raise RuntimeError(f"{path} data.list invalid")
  30. return rows
  31. def biz_key(row: dict, fields: list[str]) -> str:
  32. vals = []
  33. for f in fields:
  34. if f not in row or row[f] is None or row[f] == "":
  35. return str(row.get("bizKey") or row.get("RecID") or row.get("id") or row.get("djbh"))
  36. vals.append(str(row[f]))
  37. return "#".join(vals)
  38. def upsert(cur, table: str, source_system: str, source_table: str, row: dict, fields: list[str], batch: str, tenant_id=0):
  39. biz = biz_key(row, fields)
  40. rid = str(row.get("RecID") or row.get("id") or row.get("bizKey") or biz)
  41. raw = json.dumps(row, ensure_ascii=False)
  42. cur.execute(
  43. f"""
  44. INSERT INTO {table}
  45. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  46. raw_data, sync_batch_id, sync_time, process_status, create_time)
  47. VALUES (%s, %s, %s, %s, %s, %s, %s, NOW(), 'PENDING', NOW())
  48. ON DUPLICATE KEY UPDATE
  49. source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  50. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  51. process_status='PENDING', update_time=NOW()
  52. """,
  53. (tenant_id, source_system, source_table, rid, biz, raw, batch),
  54. )
  55. def transform_s5(cur, source_system: str, batch: str):
  56. cur.execute(
  57. """
  58. INSERT INTO mdp_std_purchase_receipt
  59. (tenant_id, factory_id, source_system, domain, receiver, line, rct_date, supp, sort_name,
  60. item_num, item_name, item_spec, um, qty_ordered, qty_received, lot_serial, location,
  61. ord_nbr, ord_line, blanket_line, pur_ord, pur_line, sales_job, address1,
  62. req, req_line, dop_req, source_biz_key, sync_batch_id, sync_time)
  63. SELECT
  64. IFNULL(p.tenant_id, 0), 1, %s,
  65. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')),
  66. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')),
  67. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')) AS SIGNED),
  68. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RctDate')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  69. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')),
  70. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')),
  71. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')), NULL, NULL,
  72. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.UM')),
  73. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)),
  74. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyReceived')) AS DECIMAL(18,6)),
  75. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.LotSerial')),
  76. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Location')),
  77. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')),
  78. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED),
  79. NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL,
  80. IFNULL(NULLIF(p.source_biz_key,''), CONCAT(
  81. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')), '#',
  82. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')), '#',
  83. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')))),
  84. %s, NOW()
  85. FROM mdp_stg_purchase_receipt p
  86. INNER JOIN mdp_stg_purchase_receipt d
  87. ON d.source_table='PurOrdRctMaster' AND d.source_system=%s AND d.sync_batch_id=%s
  88. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  89. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver'))
  90. WHERE p.source_table='PurOrdRctDetail' AND p.source_system=%s AND p.sync_batch_id=%s
  91. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc'
  92. ON DUPLICATE KEY UPDATE
  93. qty_received=VALUES(qty_received), sync_batch_id=VALUES(sync_batch_id),
  94. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  95. """,
  96. (source_system, batch, source_system, batch, source_system, batch),
  97. )
  98. def transform_s6(cur, source_system: str, batch: str):
  99. cur.execute(
  100. """
  101. INSERT INTO mdp_std_ipqc_inspection
  102. (tenant_id, factory_id, source_system, bill_no, product_model, production_batch_no, production_work_order,
  103. result_judgement, attachment, remark, inspector, process_code, process_name, production_person,
  104. sample_qty, form_no, version_no, effective_date, material_code, material_name,
  105. inspec_standard_version, inspec_standard_code, inspection_status,
  106. source_row_id, source_biz_key, sync_batch_id, sync_time)
  107. SELECT
  108. IFNULL(m.tenant_id, 1300000000001), 1, %s,
  109. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))),
  110. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.cplx')),
  111. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.scph')),
  112. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.lydjbh')),
  113. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgpd')),
  114. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.fj')),
  115. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bz')),
  116. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jyr')),
  117. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxbm')),
  118. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxmc')),
  119. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sczyry')),
  120. CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ybl')) AS DECIMAL(18,6)),
  121. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bdbh')),
  122. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bbh')),
  123. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sxrq')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  124. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlbm')),
  125. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlmc')),
  126. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbb')),
  127. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbh')),
  128. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.status')),
  129. IFNULL(m.source_row_id, JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))),
  130. IFNULL(NULLIF(m.source_biz_key,''), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh'))),
  131. %s, NOW()
  132. FROM mdp_stg_ipqc_pull m
  133. WHERE m.source_table='qms_gcjyd' AND m.source_system=%s AND m.sync_batch_id=%s
  134. ON DUPLICATE KEY UPDATE
  135. bill_no=VALUES(bill_no), sync_batch_id=VALUES(sync_batch_id),
  136. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  137. """,
  138. (source_system, batch, source_system, batch),
  139. )
  140. def transform_s7(cur, source_system: str, batch: str):
  141. cur.execute(
  142. """
  143. INSERT INTO mdp_std_production_receipt
  144. (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line,
  145. receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord,
  146. department, department_desc, applicant_name, item_num, item_name, item_spec, um,
  147. location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc,
  148. ord_nbr, source_biz_key, sync_batch_id, sync_time)
  149. SELECT
  150. IFNULL(n.tenant_id, 0), 1, %s,
  151. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')),
  152. CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED),
  153. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED),
  154. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')),
  155. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED),
  156. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  157. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  158. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  159. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')),
  160. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')),
  161. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')),
  162. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')),
  163. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  164. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  165. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')),
  166. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), NULL, NULL, NULL,
  167. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), NULL,
  168. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR),
  169. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)),
  170. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)),
  171. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), NULL,
  172. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')),
  173. CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#',
  174. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))),
  175. %s, NOW()
  176. FROM mdp_stg_production_receipt n
  177. INNER JOIN mdp_stg_production_receipt d
  178. ON d.source_table='NbrDetail' AND d.source_system=%s AND d.sync_batch_id=%s
  179. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  180. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  181. WHERE n.source_table='NbrMaster' AND n.source_system=%s AND n.sync_batch_id=%s
  182. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  183. AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
  184. ON DUPLICATE KEY UPDATE
  185. qty_rec=VALUES(qty_rec), sync_batch_id=VALUES(sync_batch_id),
  186. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  187. """,
  188. (source_system, batch, source_system, batch, source_system, batch),
  189. )
  190. def snapshot_keys(cur, table: str, batch: str) -> set[str]:
  191. cur.execute(
  192. f"SELECT source_biz_key AS k FROM {table} WHERE sync_batch_id=%s",
  193. (batch,),
  194. )
  195. return {str(r["k"]) for r in cur.fetchall() if r["k"] is not None}
  196. def compare(mod: str, keys_db: set[str], keys_api: set[str]) -> bool:
  197. only_db = sorted(keys_db - keys_api)
  198. only_api = sorted(keys_api - keys_db)
  199. print(f"[{mod}] db_keys={len(keys_db)} api_keys={len(keys_api)}")
  200. if not keys_db or not keys_api:
  201. print(f"[FAIL] {mod} 某一侧键集为空", file=sys.stderr)
  202. return False
  203. if keys_db != keys_api:
  204. print(f"[FAIL] {mod} 键集不等 only_db={only_db[:10]} only_api={only_api[:10]}", file=sys.stderr)
  205. return False
  206. print(f"[PASS] {mod} DB≡API source_biz_key")
  207. return True
  208. def dual_s5(conn) -> bool:
  209. masters = fetch("/api/receipt-master")
  210. details = fetch("/api/receipt")
  211. db_batch = f"CT_DB_S5_{TS}"
  212. api_batch = f"CT_API_S5_{TS}"
  213. with conn.cursor() as cur:
  214. for r in masters:
  215. upsert(cur, "mdp_stg_purchase_receipt", "AIDOPDEV_MYSQL", "PurOrdRctMaster", r, ["Domain", "Receiver"], db_batch)
  216. for r in details:
  217. upsert(cur, "mdp_stg_purchase_receipt", "AIDOPDEV_MYSQL", "PurOrdRctDetail", r, ["Domain", "Receiver", "Line"], db_batch)
  218. transform_s5(cur, "AIDOPDEV_MYSQL", db_batch)
  219. keys_db = snapshot_keys(cur, "mdp_std_purchase_receipt", db_batch)
  220. cur.execute("DELETE FROM mdp_std_purchase_receipt WHERE sync_batch_id=%s", (db_batch,))
  221. for r in masters:
  222. upsert(cur, "mdp_stg_purchase_receipt", "WMS_API", "PurOrdRctMaster", r, ["Domain", "Receiver"], api_batch)
  223. for r in details:
  224. upsert(cur, "mdp_stg_purchase_receipt", "WMS_API", "PurOrdRctDetail", r, ["Domain", "Receiver", "Line"], api_batch)
  225. transform_s5(cur, "WMS_API", api_batch)
  226. keys_api = snapshot_keys(cur, "mdp_std_purchase_receipt", api_batch)
  227. return compare("S5", keys_db, keys_api)
  228. def dual_s6(conn) -> bool:
  229. rows = fetch("/api/ipqc")
  230. db_batch = f"CT_DB_S6_{TS}"
  231. api_batch = f"CT_API_S6_{TS}"
  232. with conn.cursor() as cur:
  233. for r in rows:
  234. upsert(cur, "mdp_stg_ipqc_pull", "AIDOPDEV_MYSQL", "qms_gcjyd", r, ["djbh"], db_batch, tenant_id=1300000000001)
  235. transform_s6(cur, "AIDOPDEV_MYSQL", db_batch)
  236. keys_db = snapshot_keys(cur, "mdp_std_ipqc_inspection", db_batch)
  237. cur.execute("DELETE FROM mdp_std_ipqc_inspection WHERE sync_batch_id=%s", (db_batch,))
  238. for r in rows:
  239. upsert(cur, "mdp_stg_ipqc_pull", "WMS_API", "qms_gcjyd", r, ["djbh"], api_batch, tenant_id=1300000000001)
  240. transform_s6(cur, "WMS_API", api_batch)
  241. keys_api = snapshot_keys(cur, "mdp_std_ipqc_inspection", api_batch)
  242. return compare("S6", keys_db, keys_api)
  243. def dual_s7(conn) -> bool:
  244. masters = fetch("/api/production-receipt")
  245. details = fetch("/api/production-receipt-detail")
  246. db_batch = f"CT_DB_S7_{TS}"
  247. api_batch = f"CT_API_S7_{TS}"
  248. with conn.cursor() as cur:
  249. for r in masters:
  250. upsert(cur, "mdp_stg_production_receipt", "AIDOPDEV_MYSQL", "NbrMaster", r, ["Domain", "Nbr"], db_batch)
  251. for r in details:
  252. upsert(cur, "mdp_stg_production_receipt", "AIDOPDEV_MYSQL", "NbrDetail", r, ["Domain", "Nbr", "Line"], db_batch)
  253. transform_s7(cur, "AIDOPDEV_MYSQL", db_batch)
  254. keys_db = snapshot_keys(cur, "mdp_std_production_receipt", db_batch)
  255. cur.execute("DELETE FROM mdp_std_production_receipt WHERE sync_batch_id=%s", (db_batch,))
  256. for r in masters:
  257. upsert(cur, "mdp_stg_production_receipt", "WMS_API", "NbrMaster", r, ["Domain", "Nbr"], api_batch)
  258. for r in details:
  259. upsert(cur, "mdp_stg_production_receipt", "WMS_API", "NbrDetail", r, ["Domain", "Nbr", "Line"], api_batch)
  260. transform_s7(cur, "WMS_API", api_batch)
  261. keys_api = snapshot_keys(cur, "mdp_std_production_receipt", api_batch)
  262. return compare("S7", keys_db, keys_api)
  263. def main() -> int:
  264. conn = get_conn()
  265. try:
  266. ok = all([dual_s5(conn), dual_s6(conn), dual_s7(conn)])
  267. finally:
  268. conn.close()
  269. if not ok:
  270. return 1
  271. print(f"[PASS] S5/S6/S7 contract all green (ts={TS})")
  272. return 0
  273. if __name__ == "__main__":
  274. sys.exit(main())