run_contract_s6_report.py 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113
  1. """S6 报工契约:Mock /api/mes/report 以 DB/API 标签落 stg→std,比对 source_biz_key。"""
  2. from __future__ import annotations
  3. import json
  4. import sys
  5. from datetime import datetime
  6. import requests
  7. from _db import get_conn
  8. MOCK = "http://127.0.0.1:8018"
  9. TOKEN = "uat-mock-token"
  10. TS = datetime.now().strftime("%Y%m%d%H%M%S")
  11. PATH = "/api/mes/report"
  12. STG = "mdp_stg_s6_report"
  13. STD = "mdp_std_s6_report"
  14. FIELDS = ["noid", "kgdate"]
  15. def fetch() -> list[dict]:
  16. r = requests.get(MOCK.rstrip("/") + PATH, headers={"Authorization": f"Bearer {TOKEN}"}, timeout=30)
  17. r.raise_for_status()
  18. rows = r.json().get("data", {}).get("list")
  19. if not isinstance(rows, list):
  20. raise RuntimeError("invalid list")
  21. return rows
  22. def biz_key(row: dict) -> str:
  23. vals = []
  24. for f in FIELDS:
  25. alt = next((k for k in row if k.lower() == f.lower()), None)
  26. if alt is None or row[alt] in (None, ""):
  27. return str(row.get("bizKey") or row.get("Id") or row.get("id"))
  28. vals.append(str(row[alt]))
  29. return "#".join(vals)
  30. def upsert(cur, source_system: str, row: dict, batch: str):
  31. biz = biz_key(row)
  32. rid = str(row.get("Id") or row.get("id") or biz)
  33. cur.execute(
  34. f"""
  35. INSERT INTO {STG}
  36. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  37. raw_data, sync_batch_id, sync_time, process_status)
  38. VALUES (0, %s, 'Cj_Bg_Head_Rep', %s, %s, %s, %s, NOW(), 'PENDING')
  39. ON DUPLICATE KEY UPDATE
  40. raw_data=VALUES(raw_data), sync_batch_id=VALUES(sync_batch_id),
  41. sync_time=VALUES(sync_time), process_status='PENDING'
  42. """,
  43. (source_system, rid, biz, json.dumps(row, ensure_ascii=False), batch),
  44. )
  45. def transform(cur, source_system: str, batch: str):
  46. cur.execute(
  47. f"""
  48. INSERT INTO {STD}
  49. (tenant_id, factory_id, source_system, work_order_no, report_date, report_qty, ztid,
  50. source_row_id, source_biz_key, sync_batch_id, sync_time)
  51. SELECT 0, 1, %s,
  52. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.noid')),'null'), source_biz_key),
  53. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.kgdate')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  54. CAST(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sl')),'null'),'') AS DECIMAL(18,6)),
  55. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ztid')),'null'),
  56. source_row_id, source_biz_key, %s, NOW()
  57. FROM {STG} WHERE sync_batch_id=%s
  58. ON DUPLICATE KEY UPDATE
  59. work_order_no=VALUES(work_order_no), report_date=VALUES(report_date),
  60. report_qty=VALUES(report_qty), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  61. """,
  62. (source_system, batch, batch),
  63. )
  64. def keys(cur, batch: str) -> set[str]:
  65. cur.execute(f"SELECT source_biz_key AS k FROM {STD} WHERE sync_batch_id=%s", (batch,))
  66. return {str(r["k"]) for r in cur.fetchall() if r["k"] is not None}
  67. def main() -> int:
  68. rows = fetch()
  69. db_batch = f"RPT_DB_{TS}"
  70. api_batch = f"RPT_API_{TS}"
  71. conn = get_conn()
  72. try:
  73. with conn.cursor() as cur:
  74. for r in rows:
  75. upsert(cur, "T8_V5_SQLSERVER", r, db_batch)
  76. transform(cur, "T8_V5_SQLSERVER", db_batch)
  77. kdb = keys(cur, db_batch)
  78. cur.execute(f"DELETE FROM {STD} WHERE sync_batch_id=%s", (db_batch,))
  79. cur.execute(f"DELETE FROM {STG} WHERE sync_batch_id=%s", (db_batch,))
  80. for r in rows:
  81. upsert(cur, "WMS_API", r, api_batch)
  82. transform(cur, "WMS_API", api_batch)
  83. kapi = keys(cur, api_batch)
  84. finally:
  85. conn.close()
  86. print(f"[S6_REPORT] db={len(kdb)} api={len(kapi)}")
  87. if not kdb or kdb != kapi:
  88. print(f"[FAIL] mismatch db={kdb} api={kapi}", file=sys.stderr)
  89. return 1
  90. print(f"[PASS] S6 report std contract (ts={TS})")
  91. return 0
  92. if __name__ == "__main__":
  93. raise SystemExit(main())