e2e_api_to_std_s7.py 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141
  1. """S7 成品入库(NBR/WOI) 端到端:Mock API → mdp_stg_production_receipt → mdp_std_production_receipt。
  2. 用法:
  3. set MOCK_TOKEN=uat-mock-token
  4. python e2e_api_to_std_s7.py
  5. """
  6. from __future__ import annotations
  7. import json
  8. import sys
  9. from datetime import datetime
  10. import requests
  11. from _db import get_conn
  12. MOCK = "http://127.0.0.1:8018"
  13. TOKEN = "uat-mock-token"
  14. BATCH = f"CT_API_E2E_S7_{datetime.now():%Y%m%d%H%M%S}"
  15. def fetch(path: str) -> list[dict]:
  16. r = requests.get(
  17. MOCK.rstrip("/") + path,
  18. headers={"Authorization": f"Bearer {TOKEN}"},
  19. timeout=30,
  20. )
  21. r.raise_for_status()
  22. rows = r.json().get("data", {}).get("list")
  23. if not isinstance(rows, list):
  24. raise RuntimeError(f"{path} data.list invalid")
  25. return rows
  26. def upsert_stg(cur, source_table: str, row: dict, biz_expr: list[str]):
  27. vals = []
  28. for f in biz_expr:
  29. if f not in row or row[f] is None or row[f] == "":
  30. vals = None
  31. break
  32. vals.append(str(row[f]))
  33. biz = "#".join(vals) if vals else str(row.get("bizKey") or row.get("RecID") or row.get("id"))
  34. rid = str(row.get("RecID") or row.get("bizKey") or biz)
  35. raw = json.dumps(row, ensure_ascii=False)
  36. cur.execute(
  37. """
  38. INSERT INTO mdp_stg_production_receipt
  39. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  40. raw_data, sync_batch_id, sync_time, process_status, create_time)
  41. VALUES (0, 'WMS_API', %s, %s, %s, %s, %s, NOW(), 'PENDING', NOW())
  42. ON DUPLICATE KEY UPDATE
  43. source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  44. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  45. process_status='PENDING', update_time=NOW()
  46. """,
  47. (source_table, rid, biz, raw, BATCH),
  48. )
  49. def transform_std(cur):
  50. cur.execute(
  51. """
  52. INSERT INTO mdp_std_production_receipt
  53. (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line,
  54. receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord,
  55. department, department_desc, applicant_name, item_num, item_name, item_spec, um,
  56. location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc,
  57. ord_nbr, source_biz_key, sync_batch_id, sync_time)
  58. SELECT
  59. IFNULL(n.tenant_id, 0), 1, IFNULL(NULLIF(n.source_system,''), 'WMS_API'),
  60. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')),
  61. CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED),
  62. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED),
  63. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')),
  64. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED),
  65. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  66. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  67. UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')),
  68. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')),
  69. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')),
  70. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')),
  71. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')),
  72. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  73. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')),
  74. JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')),
  75. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), NULL, NULL, NULL,
  76. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), NULL,
  77. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR),
  78. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)),
  79. CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)),
  80. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), NULL,
  81. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')),
  82. CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#',
  83. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))),
  84. %s, NOW()
  85. FROM mdp_stg_production_receipt n
  86. INNER JOIN mdp_stg_production_receipt d
  87. ON d.source_table='NbrDetail'
  88. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  89. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr'))
  90. AND d.source_system='WMS_API' AND d.sync_batch_id=%s
  91. WHERE n.source_table='NbrMaster' AND n.source_system='WMS_API' AND n.sync_batch_id=%s
  92. AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI'
  93. AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1
  94. AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL
  95. ON DUPLICATE KEY UPDATE
  96. qty_rec=VALUES(qty_rec), sync_batch_id=VALUES(sync_batch_id),
  97. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  98. """,
  99. (BATCH, BATCH, BATCH),
  100. )
  101. def main() -> int:
  102. masters = fetch("/api/production-receipt")
  103. details = fetch("/api/production-receipt-detail")
  104. conn = get_conn()
  105. try:
  106. with conn.cursor() as cur:
  107. for r in masters:
  108. upsert_stg(cur, "NbrMaster", r, ["Domain", "Nbr"])
  109. for r in details:
  110. upsert_stg(cur, "NbrDetail", r, ["Domain", "Nbr", "Line"])
  111. transform_std(cur)
  112. cur.execute(
  113. "SELECT COUNT(*) AS c, COUNT(DISTINCT source_biz_key) AS k "
  114. "FROM mdp_std_production_receipt WHERE sync_batch_id=%s",
  115. (BATCH,),
  116. )
  117. row = cur.fetchone()
  118. print(f"[PASS] batch={BATCH} std rows={row['c']} keys={row['k']}")
  119. return 0 if row["c"] > 0 else 1
  120. except Exception as e:
  121. print(f"[FAIL] {e}", file=sys.stderr)
  122. return 1
  123. finally:
  124. conn.close()
  125. if __name__ == "__main__":
  126. sys.exit(main())