using Microsoft.Extensions.Logging; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Wms; /// /// 165 ASN(mdp_stg_ship_trans / ASNBOLShipperDetail)是销售发运的权威来源。 /// InvTransHist.iss-so 只进入对账表,不写 SALES_SHIP。 /// public sealed class ShipTransNeutralProjection : ITransient { /// /// 贴源筛选。mdp_stg_ship_trans 一张表承载 ASN 明细与发运计划(LinkagePlan / ShippingPlanDetail), /// 只有 ASN 明细是发货实绩;计划行没有 ShipDate / QtyShipped,放进来会变成一批 0 数量的 SALES_SHIP。 /// 同一张 ASN 还会因 Q13 的 AIDOPDEV_MYSQL → AIDOP_NATIVE 迁移在贴源留两份, /// 故聚合口径一律先按业务字段去重,不按贴源行数算。 /// private const string AsnStagingFilter = """ s.tenant_id=@tid AND s.source_table='ASNBOLShipperDetail' AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null') IS NOT NULL """; /// ASN 发货时间。eff_date / trans_time / approved_time 同源,避免三处各写一遍。 private const string ShipDateExpr = """ STR_TO_DATE(SUBSTRING_INDEX(REPLACE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShipDate')),'null'),'T',' '),'.',1),'%Y-%m-%d %H:%i:%s') """; /// ASN 明细行号。165 推 OrdLine,历史样例里也出现过 Line,两者都取不到才退 0。 private const string OrdLineExpr = """ COALESCE(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdLine')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Line')),'null'), '0') """; private readonly ISqlSugarClient _db; private readonly MdpNeutralSourceGate _gate; private readonly ILogger _logger; public ShipTransNeutralProjection(ISqlSugarClient db, MdpNeutralSourceGate gate, ILoggerFactory loggerFactory) { _db = db; _gate = gate; _logger = loggerFactory.CreateLogger(nameof(ShipTransNeutralProjection)); } public async Task ProjectAsync(long tenantId, string? batchId, CancellationToken cancellationToken = default) { cancellationToken.ThrowIfCancellationRequested(); if (!await _gate.AllowsAsync(tenantId, "INV_TRANS", "DOPDEMORQ_SQLSERVER", syncBatchId: batchId, cancellationToken: cancellationToken)) return 0; await _db.Ado.ExecuteCommandAsync( """ CREATE TABLE IF NOT EXISTS mdp_ship_recon_diff ( id BIGINT AUTO_INCREMENT PRIMARY KEY, tenant_id BIGINT NOT NULL, source_system VARCHAR(50) NOT NULL, domain VARCHAR(50) NOT NULL, ord_nbr VARCHAR(120) NOT NULL, ord_line VARCHAR(50) NOT NULL, biz_date DATE NOT NULL, asn_qty DECIMAL(18,6) NOT NULL DEFAULT 0, inv_iss_so_qty DECIMAL(18,6) NOT NULL DEFAULT 0, diff_qty DECIMAL(18,6) NOT NULL, detected_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_ship_recon (tenant_id, source_system, domain, ord_nbr, ord_line, biz_date) ) """); var inserted = await _db.Ado.ExecuteCommandAsync( $""" INSERT INTO mdp_std_inv_trans (tenant_id, source_system, written_by, domain, src_rec_id, trans_type, src_trans_type_raw, biz_doc_type, src_biz_doc_type_raw, approved_flag, void_flag, summary_flag, item_num, lot_serial, location, dimension1, dimension2, site, qty_change, doc_qty, eff_date, trans_time, approved_time, ord_nbr, ref_task_no, shipper_num, create_user, history_from, as_of, sync_batch_id) SELECT @tid, 'DOPDEMORQ_SQLSERVER', 'DB_SYNC', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),''), LEFT(CONCAT('ASN:', IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),''), s.source_row_id), ':', {OrdLineExpr}), 200), 'FG_SHIP', NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ShType')),'null'), 'SALES_SHIP', 'ASNBOLShipperDetail', IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IsConfirm')),'0') IN ('1','true','True'), 1, 0), IF(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Canceled')),'0') NOT IN ('0','') OR IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.IsActive')),'1') IN ('0','false'), 1, 0), 0, COALESCE(wol.item_code, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CustItem')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ContainerItem')),'null')), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.LotSerial')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Location')),'null'), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Dimension1')),'null'),''), IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Dimension2')),'null'),''), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Site')),'null'), -ABS(IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyShipped')),'null') AS DECIMAL(18,6)), 0)), CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyToShip')),'null') AS DECIMAL(18,6)), {ShipDateExpr}, {ShipDateExpr}, {ShipDateExpr}, NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),'null'), NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CreateUser')),'null'), CURDATE(), NOW(), IFNULL(@batch, s.sync_batch_id) FROM mdp_stg_ship_trans s LEFT JOIN mdp_std_work_order_line wol ON wol.tenant_id=@tid AND wol.source_system='AIDOP_NATIVE' AND wol.doc_type='SALES_ORDER' AND wol.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null') AND wol.line_no=CAST({OrdLineExpr} AS SIGNED) WHERE {AsnStagingFilter} ON DUPLICATE KEY UPDATE qty_change=VALUES(qty_change), doc_qty=VALUES(doc_qty), item_num=VALUES(item_num), location=VALUES(location), lot_serial=VALUES(lot_serial), trans_type=VALUES(trans_type), biz_doc_type=VALUES(biz_doc_type), eff_date=VALUES(eff_date), trans_time=VALUES(trans_time), approved_time=VALUES(approved_time), approved_flag=VALUES(approved_flag), void_flag=VALUES(void_flag), as_of=VALUES(as_of), sync_batch_id=VALUES(sync_batch_id) """, new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batchId ?? "")); try { // gate log 是诊断记录,写失败不得中断发运投影。 // 元素是 OrdNbr:OrdLine,按 70 字节估算,12 条低于 group_concat_max_len=1024。 // 外层 LEFT(...,500) 只是列宽兜底。row_count 仍是全量 COUNT(DISTINCT ...)。 await _db.Ado.ExecuteCommandAsync( $""" INSERT INTO mdp_source_gate_log (tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id) SELECT @tid, 'INV_TRANS', 'DOPDEMORQ_SQLSERVER', 'SO_LINE_NOT_FOUND', COUNT(DISTINCT CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),''), ':', {OrdLineExpr})) AS asn_rows, LEFT(( SELECT GROUP_CONCAT(x.sample_key ORDER BY x.sample_key SEPARATOR ',') FROM ( SELECT DISTINCT CONCAT(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.OrdNbr')),''), ':', {OrdLineExpr.Replace("s.raw_data", "s2.raw_data")}) AS sample_key FROM mdp_stg_ship_trans s2 LEFT JOIN mdp_std_work_order_line wol2 ON wol2.tenant_id=@tid AND wol2.source_system='AIDOP_NATIVE' AND wol2.doc_type='SALES_ORDER' AND wol2.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.OrdNbr')),'null') AND wol2.line_no=CAST({OrdLineExpr.Replace("s.raw_data", "s2.raw_data")} AS SIGNED) WHERE {AsnStagingFilter.Replace("s.tenant_id", "s2.tenant_id").Replace("s.source_table", "s2.source_table").Replace("s.raw_data", "s2.raw_data")} AND wol2.id IS NULL AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.CustItem')),'null') IS NULL AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s2.raw_data,'$.ContainerItem')),'null') IS NULL ORDER BY sample_key LIMIT 12 ) x ), 500), @batch FROM mdp_stg_ship_trans s LEFT JOIN mdp_std_work_order_line wol ON wol.tenant_id=@tid AND wol.source_system='AIDOP_NATIVE' AND wol.doc_type='SALES_ORDER' AND wol.order_no=NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null') AND wol.line_no=CAST({OrdLineExpr} AS SIGNED) WHERE {AsnStagingFilter} AND wol.id IS NULL AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.CustItem')),'null') IS NULL AND NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.ContainerItem')),'null') IS NULL HAVING asn_rows>0 """, new SugarParameter("@tid", tenantId), new SugarParameter("@batch", batchId ?? "")); } catch (Exception ex) { _logger.LogWarning(ex, "[ShipTransNeutralProjection] SO_LINE_NOT_FOUND gate log 写入失败,不影响发运投影。tenant={Tenant}", tenantId); } await _db.Ado.ExecuteCommandAsync( $""" INSERT INTO mdp_ship_recon_diff (tenant_id, source_system, domain, ord_nbr, ord_line, biz_date, asn_qty, inv_iss_so_qty, diff_qty) SELECT @tid, 'DOPDEMORQ_SQLSERVER', d.domain, d.ord_nbr, d.ord_line, d.biz_date, d.asn_qty, d.iss_qty, d.asn_qty - d.iss_qty FROM ( SELECT domain, ord_nbr, IFNULL(ord_line,'') ord_line, biz_date, SUM(asn_qty) asn_qty, SUM(iss_qty) iss_qty FROM ( SELECT a.domain, a.ord_nbr, a.ord_line, a.biz_date, a.asn_qty, 0 iss_qty FROM ( SELECT IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Domain')),'null'),'') domain, IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdNbr')),'null'),'') ord_nbr, IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.OrdLine')),'null'),'') ord_line, DATE({ShipDateExpr}) biz_date, IFNULL(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.Id')),'') asn_id, MAX(ABS(IFNULL(CAST(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(s.raw_data,'$.QtyShipped')),'null') AS DECIMAL(18,6)),0))) asn_qty FROM mdp_stg_ship_trans s WHERE {AsnStagingFilter} GROUP BY domain, ord_nbr, ord_line, biz_date, asn_id ) a UNION ALL SELECT IFNULL(t.domain,''), IFNULL(t.ord_nbr,''), '', DATE(t.eff_date), 0, ABS(IFNULL(t.qty_change,0)) FROM mdp_std_inv_trans t WHERE t.tenant_id=@tid AND t.source_system='DOPDEMORQ_SQLSERVER' AND t.src_trans_type_raw='iss-so' AND t.eff_date>=DATE_SUB(CURDATE(), INTERVAL 14 DAY) ) u WHERE biz_date IS NOT NULL AND ord_nbr<>'' GROUP BY domain, ord_nbr, ord_line, biz_date ) d WHERE ABS(d.asn_qty - d.iss_qty) > 0.0001 ON DUPLICATE KEY UPDATE asn_qty=VALUES(asn_qty), inv_iss_so_qty=VALUES(inv_iss_so_qty), diff_qty=VALUES(diff_qty), detected_at=CURRENT_TIMESTAMP """, new SugarParameter("@tid", tenantId)); return inserted; } }