ShipTransNeutralProjection.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214
  1. using Microsoft.Extensions.Logging;
  2. using SqlSugar;
  3. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Wms;
  4. /// <summary>
  5. /// 165 ASN(mdp_stg_ship_trans / ASNBOLShipperDetail)是销售发运的权威来源。
  6. /// InvTransHist.iss-so 只进入对账表,不写 SALES_SHIP。
  7. /// </summary>
  8. public sealed class ShipTransNeutralProjection : ITransient
  9. {
  10. /// <summary>
  11. /// 贴源筛选。<c>mdp_stg_ship_trans</c> 一张表承载 ASN 明细与发运计划(LinkagePlan / ShippingPlanDetail),
  12. /// 只有 ASN 明细是发货实绩;计划行没有 ShipDate / QtyShipped,放进来会变成一批 0 数量的 SALES_SHIP。
  13. /// 同一张 ASN 还会因 Q13 的 <c>AIDOPDEV_MYSQL → AIDOP_NATIVE</c> 迁移在贴源留两份,
  14. /// 故聚合口径一律先按业务字段去重,不按贴源行数算。
  15. /// </summary>
  16. private const string AsnStagingFilter = """
  17. s.tenant_id=@tid
  18. AND s.source_table='ASNBOLShipperDetail'
  19. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null') IS NOT NULL
  20. """;
  21. /// <summary>ASN 发货时间。eff_date / trans_time / approved_time 同源,避免三处各写一遍。</summary>
  22. private const string ShipDateExpr = """
  23. STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s')
  24. """;
  25. /// <summary>ASN 明细行号。165 推 OrdLine,历史样例里也出现过 Line,两者都取不到才退 0。</summary>
  26. private const string OrdLineExpr = """
  27. COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdLine')),'null'),
  28. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Line')),'null'), '0')
  29. """;
  30. private readonly ISqlSugarClient _db;
  31. private readonly MdpNeutralSourceGate _gate;
  32. private readonly ILogger _logger;
  33. public ShipTransNeutralProjection(ISqlSugarClient db, MdpNeutralSourceGate gate, ILoggerFactory loggerFactory)
  34. {
  35. _db = db;
  36. _gate = gate;
  37. _logger = loggerFactory.CreateLogger(nameof(ShipTransNeutralProjection));
  38. }
  39. public async Task<int> ProjectAsync(long tenantId, string? batchId, CancellationToken cancellationToken = default)
  40. {
  41. cancellationToken.ThrowIfCancellationRequested();
  42. if (!await _gate.AllowsAsync(tenantId, "INV_TRANS", "DOPDEMORQ_SQLSERVER", syncBatchId: batchId, cancellationToken: cancellationToken))
  43. return 0;
  44. await _db.Ado.ExecuteCommandAsync(
  45. """
  46. CREATE TABLE IF NOT EXISTS mdp_ship_recon_diff (
  47. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  48. tenant_id BIGINT NOT NULL,
  49. source_system VARCHAR(50) NOT NULL,
  50. domain VARCHAR(50) NOT NULL,
  51. ord_nbr VARCHAR(120) NOT NULL,
  52. ord_line VARCHAR(50) NOT NULL,
  53. biz_date DATE NOT NULL,
  54. asn_qty DECIMAL(18,6) NOT NULL DEFAULT 0,
  55. inv_iss_so_qty DECIMAL(18,6) NOT NULL DEFAULT 0,
  56. diff_qty DECIMAL(18,6) NOT NULL,
  57. detected_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  58. UNIQUE KEY uk_ship_recon (tenant_id, source_system, domain, ord_nbr, ord_line, biz_date)
  59. )
  60. """);
  61. var inserted = await _db.Ado.ExecuteCommandAsync(
  62. $"""
  63. INSERT INTO mdp_std_inv_trans
  64. (tenant_id, source_system, written_by, domain, src_rec_id, trans_type, src_trans_type_raw,
  65. biz_doc_type, src_biz_doc_type_raw, approved_flag, void_flag, summary_flag,
  66. item_num, lot_serial, location, dimension1, dimension2, site, qty_change, doc_qty,
  67. eff_date, trans_time, approved_time, ord_nbr, ref_task_no, shipper_num, create_user,
  68. history_from, as_of, sync_batch_id)
  69. SELECT
  70. @tid, 'DOPDEMORQ_SQLSERVER', 'DB_SYNC',
  71. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''),
  72. LEFT(CONCAT('ASN:', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),''), s.source_row_id), ':', {OrdLineExpr}), 200),
  73. 'FG_SHIP',
  74. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShType')),'null'),
  75. 'SALES_SHIP', 'ASNBOLShipperDetail',
  76. IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IsConfirm')),'0') IN ('1','true','True'), 1, 0),
  77. IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Canceled')),'0') NOT IN ('0','')
  78. OR IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IsActive')),'1') IN ('0','false'), 1, 0),
  79. 0,
  80. COALESCE(wol.item_code,
  81. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CustItem')),'null'),
  82. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ContainerItem')),'null')),
  83. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LotSerial')),'null'),
  84. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Location')),'null'),
  85. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Dimension1')),'null'),''),
  86. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Dimension2')),'null'),''),
  87. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Site')),'null'),
  88. -ABS(IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyShipped')),'null') AS DECIMAL(18,6)), 0)),
  89. CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyToShip')),'null') AS DECIMAL(18,6)),
  90. {ShipDateExpr},
  91. {ShipDateExpr},
  92. {ShipDateExpr},
  93. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'),
  94. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'),
  95. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),'null'),
  96. NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CreateUser')),'null'),
  97. CURDATE(), NOW(), IFNULL(@batch, s.sync_batch_id)
  98. FROM mdp_stg_ship_trans s
  99. LEFT JOIN mdp_std_work_order_line wol
  100. ON wol.tenant_id=@tid AND wol.source_system='AIDOP_NATIVE' AND wol.doc_type='SALES_ORDER'
  101. AND wol.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null')
  102. AND wol.line_no=CAST({OrdLineExpr} AS SIGNED)
  103. WHERE {AsnStagingFilter}
  104. ON DUPLICATE KEY UPDATE
  105. qty_change=VALUES(qty_change), doc_qty=VALUES(doc_qty), item_num=VALUES(item_num),
  106. location=VALUES(location), lot_serial=VALUES(lot_serial), trans_type=VALUES(trans_type),
  107. biz_doc_type=VALUES(biz_doc_type), eff_date=VALUES(eff_date),
  108. trans_time=VALUES(trans_time), approved_time=VALUES(approved_time),
  109. approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag),
  110. as_of=VALUES(as_of), sync_batch_id=VALUES(sync_batch_id)
  111. """,
  112. new SugarParameter("@tid", tenantId),
  113. new SugarParameter("@batch", batchId ?? ""));
  114. try
  115. {
  116. // gate log 是诊断记录,写失败不得中断发运投影。
  117. // 元素是 OrdNbr:OrdLine,按 70 字节估算,12 条低于 group_concat_max_len=1024。
  118. // 外层 LEFT(...,500) 只是列宽兜底。row_count 仍是全量 COUNT(DISTINCT ...)。
  119. await _db.Ado.ExecuteCommandAsync(
  120. $"""
  121. INSERT INTO mdp_source_gate_log
  122. (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
  123. SELECT @tid, 'INV_TRANS', 'DOPDEMORQ_SQLSERVER', 'SO_LINE_NOT_FOUND',
  124. COUNT(DISTINCT CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),''), ':', {OrdLineExpr})) AS asn_rows,
  125. LEFT((
  126. SELECT GROUP_CONCAT(x.sample_key ORDER BY x.sample_key SEPARATOR ',')
  127. FROM (
  128. SELECT DISTINCT CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.OrdNbr')),''), ':', {OrdLineExpr.Replace("s.raw_data", "s2.raw_data")}) AS sample_key
  129. FROM mdp_stg_ship_trans s2
  130. LEFT JOIN mdp_std_work_order_line wol2
  131. ON wol2.tenant_id=@tid AND wol2.source_system='AIDOP_NATIVE' AND wol2.doc_type='SALES_ORDER'
  132. AND wol2.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.OrdNbr')),'null')
  133. AND wol2.line_no=CAST({OrdLineExpr.Replace("s.raw_data", "s2.raw_data")} AS SIGNED)
  134. WHERE {AsnStagingFilter.Replace("s.tenant_id", "s2.tenant_id").Replace("s.source_table", "s2.source_table").Replace("s.raw_data", "s2.raw_data")}
  135. AND wol2.id IS NULL
  136. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.CustItem')),'null') IS NULL
  137. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.ContainerItem')),'null') IS NULL
  138. ORDER BY sample_key
  139. LIMIT 12
  140. ) x
  141. ), 500),
  142. @batch
  143. FROM mdp_stg_ship_trans s
  144. LEFT JOIN mdp_std_work_order_line wol
  145. ON wol.tenant_id=@tid AND wol.source_system='AIDOP_NATIVE' AND wol.doc_type='SALES_ORDER'
  146. AND wol.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null')
  147. AND wol.line_no=CAST({OrdLineExpr} AS SIGNED)
  148. WHERE {AsnStagingFilter}
  149. AND wol.id IS NULL
  150. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CustItem')),'null') IS NULL
  151. AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ContainerItem')),'null') IS NULL
  152. HAVING asn_rows>0
  153. """,
  154. new SugarParameter("@tid", tenantId),
  155. new SugarParameter("@batch", batchId ?? ""));
  156. }
  157. catch (Exception ex)
  158. {
  159. _logger.LogWarning(ex,
  160. "[ShipTransNeutralProjection] SO_LINE_NOT_FOUND gate log 写入失败,不影响发运投影。tenant={Tenant}",
  161. tenantId);
  162. }
  163. await _db.Ado.ExecuteCommandAsync(
  164. $"""
  165. INSERT INTO mdp_ship_recon_diff
  166. (tenant_id, source_system, domain, ord_nbr, ord_line, biz_date, asn_qty, inv_iss_so_qty, diff_qty)
  167. SELECT @tid, 'DOPDEMORQ_SQLSERVER', d.domain, d.ord_nbr, d.ord_line, d.biz_date,
  168. d.asn_qty, d.iss_qty, d.asn_qty - d.iss_qty
  169. FROM (
  170. SELECT domain, ord_nbr, IFNULL(ord_line,'') ord_line, biz_date,
  171. SUM(asn_qty) asn_qty, SUM(iss_qty) iss_qty
  172. FROM (
  173. SELECT a.domain, a.ord_nbr, a.ord_line, a.biz_date, a.asn_qty, 0 iss_qty
  174. FROM (
  175. SELECT IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),'') domain,
  176. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'),'') ord_nbr,
  177. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdLine')),'null'),'') ord_line,
  178. DATE({ShipDateExpr}) biz_date,
  179. IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),'') asn_id,
  180. MAX(ABS(IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyShipped')),'null') AS DECIMAL(18,6)),0))) asn_qty
  181. FROM mdp_stg_ship_trans s
  182. WHERE {AsnStagingFilter}
  183. GROUP BY domain, ord_nbr, ord_line, biz_date, asn_id
  184. ) a
  185. UNION ALL
  186. SELECT IFNULL(t.domain,''), IFNULL(t.ord_nbr,''), '', DATE(t.eff_date), 0, ABS(IFNULL(t.qty_change,0))
  187. FROM mdp_std_inv_trans t
  188. WHERE t.tenant_id=@tid AND t.source_system='DOPDEMORQ_SQLSERVER'
  189. AND t.src_trans_type_raw='iss-so'
  190. AND t.eff_date>=DATE_SUB(CURDATE(), INTERVAL 14 DAY)
  191. ) u
  192. WHERE biz_date IS NOT NULL AND ord_nbr<>''
  193. GROUP BY domain, ord_nbr, ord_line, biz_date
  194. ) d
  195. WHERE ABS(d.asn_qty - d.iss_qty) > 0.0001
  196. ON DUPLICATE KEY UPDATE
  197. asn_qty=VALUES(asn_qty), inv_iss_so_qty=VALUES(inv_iss_so_qty),
  198. diff_qty=VALUES(diff_qty), detected_at=CURRENT_TIMESTAMP
  199. """,
  200. new SugarParameter("@tid", tenantId));
  201. return inserted;
  202. }
  203. }