e2e_api_to_std_s5.py 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136
  1. """S5 端到端:Mock API → stg(Writer 同构 upsert)→ std → 校验 source_biz_key。
  2. 不依赖后端进程;直连 Mock + MySQL(环境变量同 _db.py)。
  3. 用法:
  4. set MOCK_TOKEN=uat-mock-token
  5. python e2e_api_to_std_s5.py
  6. """
  7. from __future__ import annotations
  8. import json
  9. import sys
  10. from datetime import datetime
  11. import requests
  12. from _db import get_conn
  13. MOCK = "http://127.0.0.1:8018"
  14. TOKEN = "uat-mock-token"
  15. BATCH = f"CT_API_E2E_{datetime.now():%Y%m%d%H%M%S}"
  16. def fetch(path: str) -> list[dict]:
  17. r = requests.get(
  18. MOCK.rstrip("/") + path,
  19. headers={"Authorization": f"Bearer {TOKEN}"},
  20. timeout=30,
  21. )
  22. r.raise_for_status()
  23. rows = r.json().get("data", {}).get("list")
  24. if not isinstance(rows, list):
  25. raise RuntimeError(f"{path} data.list invalid")
  26. return rows
  27. def upsert_stg(cur, source_table: str, row: dict, biz_expr: list[str]):
  28. vals = []
  29. for f in biz_expr:
  30. if f not in row or row[f] is None or row[f] == "":
  31. vals = None
  32. break
  33. vals.append(str(row[f]))
  34. biz = "#".join(vals) if vals else str(row.get("bizKey") or row.get("RecID") or row.get("id"))
  35. rid = str(row.get("bizKey") or row.get("RecID") or row.get("Line") or biz)
  36. raw = json.dumps(row, ensure_ascii=False)
  37. cur.execute(
  38. """
  39. INSERT INTO mdp_stg_purchase_receipt
  40. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  41. raw_data, sync_batch_id, sync_time, process_status, create_time)
  42. VALUES (0, 'WMS_API', %s, %s, %s, %s, %s, NOW(), 'PENDING', NOW())
  43. ON DUPLICATE KEY UPDATE
  44. source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  45. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  46. process_status='PENDING', update_time=NOW()
  47. """,
  48. (source_table, rid, biz, raw, BATCH),
  49. )
  50. def transform_std(cur):
  51. cur.execute(
  52. """
  53. INSERT INTO mdp_std_purchase_receipt
  54. (tenant_id, factory_id, source_system, domain, receiver, line, rct_date, supp, sort_name,
  55. item_num, item_name, item_spec, um, qty_ordered, qty_received, lot_serial, location,
  56. ord_nbr, ord_line, blanket_line, pur_ord, pur_line, sales_job, address1,
  57. req, req_line, dop_req, source_biz_key, sync_batch_id, sync_time)
  58. SELECT
  59. IFNULL(p.tenant_id, 0), 1, IFNULL(NULLIF(p.source_system,''), 'WMS_API'),
  60. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')),
  61. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')),
  62. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')) AS SIGNED),
  63. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RctDate')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  64. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Supp')),
  65. JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Supp')),
  66. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.ItemNum')), NULL, NULL,
  67. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.UM')),
  68. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyOrded')) AS DECIMAL(18,6)),
  69. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.QtyReceived')) AS DECIMAL(18,6)),
  70. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.LotSerial')),
  71. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Location')),
  72. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdNbr')),
  73. CAST(JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.OrdLine')) AS SIGNED),
  74. NULL, NULL, NULL, NULL, NULL,
  75. NULL, NULL, NULL,
  76. IFNULL(NULLIF(p.source_biz_key,''), CONCAT(
  77. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')), '#',
  78. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')), '#',
  79. JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Line')))),
  80. %s, NOW()
  81. FROM mdp_stg_purchase_receipt p
  82. INNER JOIN mdp_stg_purchase_receipt d
  83. ON d.source_table='PurOrdRctMaster'
  84. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain'))
  85. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.Receiver')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Receiver'))
  86. AND d.source_system='WMS_API'
  87. WHERE p.source_table='PurOrdRctDetail' AND p.source_system='WMS_API'
  88. AND JSON_UNQUOTE(JSON_EXTRACT(p.raw_data,'$.RctType'))='rc'
  89. AND p.sync_batch_id=%s
  90. ON DUPLICATE KEY UPDATE
  91. qty_received=VALUES(qty_received), sync_batch_id=VALUES(sync_batch_id),
  92. sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  93. """,
  94. (BATCH, BATCH),
  95. )
  96. def main() -> int:
  97. details = fetch("/api/receipt")
  98. masters = fetch("/api/receipt-master")
  99. conn = get_conn()
  100. try:
  101. with conn.cursor() as cur:
  102. for r in masters:
  103. upsert_stg(cur, "PurOrdRctMaster", r, ["Domain", "Receiver"])
  104. for r in details:
  105. upsert_stg(cur, "PurOrdRctDetail", r, ["Domain", "Receiver", "Line"])
  106. transform_std(cur)
  107. cur.execute(
  108. "SELECT COUNT(*) AS c, COUNT(DISTINCT source_biz_key) AS k "
  109. "FROM mdp_std_purchase_receipt WHERE sync_batch_id=%s",
  110. (BATCH,),
  111. )
  112. row = cur.fetchone()
  113. print(f"[PASS] batch={BATCH} std rows={row['c']} keys={row['k']}")
  114. return 0 if row["c"] > 0 else 1
  115. except Exception as e:
  116. print(f"[FAIL] {e}", file=sys.stderr)
  117. return 1
  118. finally:
  119. conn.close()
  120. if __name__ == "__main__":
  121. sys.exit(main())