e2e_api_to_std_s6.py 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122
  1. """S6 IPQC 端到端:Mock API → mdp_stg_ipqc_pull → mdp_std_ipqc_inspection。
  2. 用法:
  3. set MOCK_TOKEN=uat-mock-token
  4. python e2e_api_to_std_s6.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_S6_{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, row: dict):
  27. biz = str(row.get("djbh") or row.get("bizKey") or row.get("id"))
  28. rid = str(row.get("id") or biz)
  29. raw = json.dumps(row, ensure_ascii=False)
  30. cur.execute(
  31. """
  32. INSERT INTO mdp_stg_ipqc_pull
  33. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  34. raw_data, sync_batch_id, sync_time, process_status, create_time)
  35. VALUES (1300000000001, 'WMS_API', 'qms_gcjyd', %s, %s, %s, %s, NOW(), 'PENDING', NOW())
  36. ON DUPLICATE KEY UPDATE
  37. source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  38. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  39. process_status='PENDING', update_time=NOW()
  40. """,
  41. (rid, biz, raw, BATCH),
  42. )
  43. def transform_std(cur):
  44. cur.execute(
  45. """
  46. INSERT INTO mdp_std_ipqc_inspection
  47. (tenant_id, factory_id, source_system, bill_no, product_model, production_batch_no, production_work_order,
  48. result_judgement, attachment, remark, inspector, process_code, process_name, production_person,
  49. sample_qty, form_no, version_no, effective_date, material_code, material_name,
  50. inspec_standard_version, inspec_standard_code, inspection_status,
  51. source_row_id, source_biz_key, sync_batch_id, sync_time)
  52. SELECT
  53. IFNULL(m.tenant_id, 1300000000001), 1, IFNULL(NULLIF(m.source_system,''), 'WMS_API'),
  54. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))),
  55. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.cplx')),
  56. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.scph')),
  57. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.lydjbh')),
  58. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgpd')),
  59. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.fj')),
  60. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bz')),
  61. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jyr')),
  62. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxbm')),
  63. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxmc')),
  64. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sczyry')),
  65. CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ybl')) AS DECIMAL(18,6)),
  66. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bdbh')),
  67. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bbh')),
  68. STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sxrq')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
  69. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlbm')),
  70. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlmc')),
  71. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbb')),
  72. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbh')),
  73. JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.status')),
  74. IFNULL(m.source_row_id, JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))),
  75. IFNULL(NULLIF(m.source_biz_key,''), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh'))),
  76. %s, NOW()
  77. FROM mdp_stg_ipqc_pull m
  78. WHERE m.source_table='qms_gcjyd' AND m.source_system='WMS_API' AND m.sync_batch_id=%s
  79. ON DUPLICATE KEY UPDATE
  80. bill_no=VALUES(bill_no), result_judgement=VALUES(result_judgement),
  81. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP
  82. """,
  83. (BATCH, BATCH),
  84. )
  85. def main() -> int:
  86. rows = fetch("/api/ipqc")
  87. conn = get_conn()
  88. try:
  89. with conn.cursor() as cur:
  90. for r in rows:
  91. upsert_stg(cur, r)
  92. transform_std(cur)
  93. cur.execute(
  94. "SELECT COUNT(*) AS c, COUNT(DISTINCT source_biz_key) AS k "
  95. "FROM mdp_std_ipqc_inspection WHERE sync_batch_id=%s",
  96. (BATCH,),
  97. )
  98. row = cur.fetchone()
  99. print(f"[PASS] batch={BATCH} std rows={row['c']} keys={row['k']}")
  100. return 0 if row["c"] > 0 else 1
  101. except Exception as e:
  102. print(f"[FAIL] {e}", file=sys.stderr)
  103. return 1
  104. finally:
  105. conn.close()
  106. if __name__ == "__main__":
  107. sys.exit(main())