|
|
@@ -1,4 +1,5 @@
|
|
|
using System.Collections.Concurrent;
|
|
|
+using Admin.NET.Plugin.AiDOP.DataPlatform;
|
|
|
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
|
|
|
using Microsoft.Extensions.Logging;
|
|
|
using SqlSugar;
|
|
|
@@ -17,9 +18,27 @@ public sealed class MdpStagingWriter : ITransient
|
|
|
"raw_data", "sync_batch_id", "sync_time", "process_status", "create_time", "update_time"
|
|
|
];
|
|
|
|
|
|
+ /// <summary>
|
|
|
+ /// 单条多值 INSERT 的最大行数。每行约 12 个占位符,200 行约 2400 个,
|
|
|
+ /// 远低于 MySQL 的 65535 上限,同时让单条语句的 redo 量保持可控。
|
|
|
+ /// </summary>
|
|
|
+ internal const int UpsertBatchRows = 200;
|
|
|
+
|
|
|
private readonly ISqlSugarClient _db;
|
|
|
private readonly ILogger<MdpStagingWriter> _logger;
|
|
|
|
|
|
+ /// <summary>一行已解析好、可直接拼进 SQL 的贴源值。</summary>
|
|
|
+ internal sealed record MdpStagingUpsertRow(
|
|
|
+ long TenantId,
|
|
|
+ long? FactoryId,
|
|
|
+ long? CompanyId,
|
|
|
+ string SourceSystem,
|
|
|
+ string SourceTable,
|
|
|
+ string SourceRowId,
|
|
|
+ string BizKey,
|
|
|
+ string RawJson,
|
|
|
+ string BatchId);
|
|
|
+
|
|
|
public MdpStagingWriter(ISqlSugarClient db, ILogger<MdpStagingWriter> logger)
|
|
|
{
|
|
|
_db = db;
|
|
|
@@ -78,47 +97,148 @@ public sealed class MdpStagingWriter : ITransient
|
|
|
}
|
|
|
var companyValue = TryParseOptionalLong(row, "company_id");
|
|
|
|
|
|
- var insertCols = new List<string>();
|
|
|
- var insertVals = new List<string>();
|
|
|
- var parameters = new List<SugarParameter>();
|
|
|
+ var resolved = new MdpStagingUpsertRow(
|
|
|
+ tenantValue.Value,
|
|
|
+ factoryValue,
|
|
|
+ companyValue,
|
|
|
+ MdpSourceIdentity.NeutralCode(source.SystemCode, source.SourceCode),
|
|
|
+ sourceTable,
|
|
|
+ sourceRowId,
|
|
|
+ bizKey,
|
|
|
+ rawJson,
|
|
|
+ ctx.BatchId);
|
|
|
+
|
|
|
+ var (sql, parameters) = BuildUpsertSql(table, columns, [resolved], now);
|
|
|
+ return await _db.Ado.ExecuteCommandAsync(sql, parameters);
|
|
|
+ }
|
|
|
|
|
|
- if (columns.Contains("tenant_id"))
|
|
|
- {
|
|
|
- insertCols.Add("tenant_id");
|
|
|
- insertVals.Add("@tenant");
|
|
|
- parameters.Add(new SugarParameter("@tenant", tenantValue.Value));
|
|
|
- }
|
|
|
+ /// <summary>
|
|
|
+ /// 同表、同实体、同批次的多行一次写入。返回与入参等长的逐行结果(0 表示租户无法解析或作用域不匹配已跳过)。
|
|
|
+ /// 超过 <see cref="UpsertBatchRows"/> 时切成多条语句,不把整批拼成一条。
|
|
|
+ /// </summary>
|
|
|
+ public async Task<IReadOnlyList<int>> UpsertBatchAsync(
|
|
|
+ MdpSource source,
|
|
|
+ MdpEntity entity,
|
|
|
+ string sourceTable,
|
|
|
+ IReadOnlyList<(IDictionary<string, object?> Row, string RawJson, string SourceRowId)> rows,
|
|
|
+ MdpPullContext ctx)
|
|
|
+ {
|
|
|
+ var results = new int[rows.Count];
|
|
|
+ if (rows.Count == 0)
|
|
|
+ return results;
|
|
|
+
|
|
|
+ if (string.IsNullOrWhiteSpace(entity.TargetTableName))
|
|
|
+ throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 target_table_name");
|
|
|
|
|
|
- if (columns.Contains("factory_id"))
|
|
|
+ var table = entity.TargetTableName!.Trim();
|
|
|
+ if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
|
|
|
+ throw new InvalidOperationException($"非法 target_table_name:{table}");
|
|
|
+
|
|
|
+ var columns = await GetTableColumnsAsync(table);
|
|
|
+ ValidateTableContract(table, entity.EntityCode, columns);
|
|
|
+ var now = DateTime.Now;
|
|
|
+ var sourceSystem = MdpSourceIdentity.NeutralCode(source.SystemCode, source.SourceCode);
|
|
|
+
|
|
|
+ var pending = new List<(int Index, MdpStagingUpsertRow Row)>();
|
|
|
+ for (var i = 0; i < rows.Count; i++)
|
|
|
{
|
|
|
- insertCols.Add("factory_id");
|
|
|
- insertVals.Add("@factory");
|
|
|
- parameters.Add(new SugarParameter("@factory", factoryValue.HasValue ? factoryValue.Value : DBNull.Value));
|
|
|
+ var (row, rawJson, sourceRowId) = rows[i];
|
|
|
+ var tenantValue = ResolveTenantId(row, ctx);
|
|
|
+ if (tenantValue is null)
|
|
|
+ {
|
|
|
+ _logger.LogWarning(
|
|
|
+ "贴源写入跳过:无法解析 tenant_id,entity={EntityCode}, source={SourceCode}, table={Table}, rowId={RowId}",
|
|
|
+ entity.EntityCode, source.SourceCode, table, sourceRowId);
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ if (ctx.TenantId > 0 && tenantValue.Value != ctx.TenantId)
|
|
|
+ continue;
|
|
|
+
|
|
|
+ var factoryValue = NormalizeFactoryId(
|
|
|
+ tenantValue.Value,
|
|
|
+ TryParseOptionalLong(row, "factory_id") ?? TryParseOptionalLong(row, "FactoryId"));
|
|
|
+ if (ctx.FactoryId > 0)
|
|
|
+ {
|
|
|
+ var resolvedFactory = factoryValue is > 0 ? factoryValue.Value : 1;
|
|
|
+ if (resolvedFactory != ctx.FactoryId)
|
|
|
+ continue;
|
|
|
+ factoryValue = resolvedFactory;
|
|
|
+ }
|
|
|
+
|
|
|
+ pending.Add((i, new MdpStagingUpsertRow(
|
|
|
+ tenantValue.Value,
|
|
|
+ factoryValue,
|
|
|
+ TryParseOptionalLong(row, "company_id"),
|
|
|
+ sourceSystem,
|
|
|
+ sourceTable,
|
|
|
+ sourceRowId,
|
|
|
+ BuildBizKey(entity.BizKeyExpr, row) ?? sourceRowId,
|
|
|
+ rawJson,
|
|
|
+ ctx.BatchId)));
|
|
|
}
|
|
|
|
|
|
- if (columns.Contains("company_id"))
|
|
|
+ for (var offset = 0; offset < pending.Count; offset += UpsertBatchRows)
|
|
|
{
|
|
|
- insertCols.Add("company_id");
|
|
|
- insertVals.Add("@company");
|
|
|
- parameters.Add(new SugarParameter("@company", companyValue.HasValue ? companyValue.Value : DBNull.Value));
|
|
|
+ var chunk = pending.Skip(offset).Take(UpsertBatchRows).ToList();
|
|
|
+ var (sql, parameters) = BuildUpsertSql(table, columns, chunk.Select(x => x.Row).ToList(), now);
|
|
|
+ await _db.Ado.ExecuteCommandAsync(sql, parameters);
|
|
|
+ foreach (var (index, _) in chunk)
|
|
|
+ results[index] = 1;
|
|
|
}
|
|
|
|
|
|
+ return results;
|
|
|
+ }
|
|
|
+
|
|
|
+ /// <summary>
|
|
|
+ /// 拼一条多值 INSERT ... ON DUPLICATE KEY UPDATE。
|
|
|
+ /// 参数名是「基名 + 下划线 + 行号」。下划线不可省:行号是纯数字,直接拼接会让
|
|
|
+ /// <c>@now</c> 的第 20 行与 <c>@now2</c> 的第 0 行都得到 <c>@now20</c>,
|
|
|
+ /// 凑满 21 行的批次就抛「Parameter '@now20' has already been defined」。
|
|
|
+ /// 加了分隔符后,按最后一个下划线切分即可唯一还原 (基名, 行号),基名含下划线也不冲突。
|
|
|
+ /// </summary>
|
|
|
+ internal static (string Sql, List<SugarParameter> Parameters) BuildUpsertSql(
|
|
|
+ string table,
|
|
|
+ IReadOnlySet<string> columns,
|
|
|
+ IReadOnlyList<MdpStagingUpsertRow> rows,
|
|
|
+ DateTime now)
|
|
|
+ {
|
|
|
+ var insertCols = new List<string>();
|
|
|
+ if (columns.Contains("tenant_id")) insertCols.Add("tenant_id");
|
|
|
+ if (columns.Contains("factory_id")) insertCols.Add("factory_id");
|
|
|
+ if (columns.Contains("company_id")) insertCols.Add("company_id");
|
|
|
insertCols.AddRange(
|
|
|
[
|
|
|
"source_system", "source_table", "source_row_id", "source_biz_key",
|
|
|
"raw_data", "sync_batch_id", "sync_time", "process_status", "create_time"
|
|
|
]);
|
|
|
- insertVals.AddRange(["@sys", "@tbl", "@rid", "@biz", "@raw", "@batch", "@now", "'PENDING'", "@now"]);
|
|
|
- parameters.AddRange(
|
|
|
- [
|
|
|
- new SugarParameter("@sys", source.SourceCode),
|
|
|
- new SugarParameter("@tbl", sourceTable),
|
|
|
- new SugarParameter("@rid", sourceRowId),
|
|
|
- new SugarParameter("@biz", bizKey),
|
|
|
- new SugarParameter("@raw", rawJson),
|
|
|
- new SugarParameter("@batch", ctx.BatchId),
|
|
|
- new SugarParameter("@now", now)
|
|
|
- ]);
|
|
|
+
|
|
|
+ var parameters = new List<SugarParameter>();
|
|
|
+ var valueGroups = new List<string>();
|
|
|
+ for (var i = 0; i < rows.Count; i++)
|
|
|
+ {
|
|
|
+ var row = rows[i];
|
|
|
+ var vals = new List<string>();
|
|
|
+ void Add(string name, object? value)
|
|
|
+ {
|
|
|
+ var p = $"{name}_{i}";
|
|
|
+ vals.Add(p);
|
|
|
+ parameters.Add(new SugarParameter(p, value ?? DBNull.Value));
|
|
|
+ }
|
|
|
+
|
|
|
+ if (columns.Contains("tenant_id")) Add("@tenant", row.TenantId);
|
|
|
+ if (columns.Contains("factory_id")) Add("@factory", row.FactoryId);
|
|
|
+ if (columns.Contains("company_id")) Add("@company", row.CompanyId);
|
|
|
+ Add("@sys", row.SourceSystem);
|
|
|
+ Add("@tbl", row.SourceTable);
|
|
|
+ Add("@rid", row.SourceRowId);
|
|
|
+ Add("@biz", row.BizKey);
|
|
|
+ Add("@raw", row.RawJson);
|
|
|
+ Add("@batch", row.BatchId);
|
|
|
+ Add("@now", now);
|
|
|
+ vals.Add("'PENDING'");
|
|
|
+ Add("@now2", now);
|
|
|
+ valueGroups.Add($"({string.Join(", ", vals)})");
|
|
|
+ }
|
|
|
|
|
|
var updateParts = new List<string>
|
|
|
{
|
|
|
@@ -128,26 +248,20 @@ public sealed class MdpStagingWriter : ITransient
|
|
|
"sync_time=VALUES(sync_time)",
|
|
|
"process_status='PENDING'"
|
|
|
};
|
|
|
-
|
|
|
- if (columns.Contains("tenant_id"))
|
|
|
- updateParts.Add("tenant_id=VALUES(tenant_id)");
|
|
|
- if (columns.Contains("factory_id"))
|
|
|
- updateParts.Add("factory_id=VALUES(factory_id)");
|
|
|
- if (columns.Contains("company_id"))
|
|
|
- updateParts.Add("company_id=VALUES(company_id)");
|
|
|
-
|
|
|
- updateParts.Add("update_time=@now");
|
|
|
+ if (columns.Contains("tenant_id")) updateParts.Add("tenant_id=VALUES(tenant_id)");
|
|
|
+ if (columns.Contains("factory_id")) updateParts.Add("factory_id=VALUES(factory_id)");
|
|
|
+ if (columns.Contains("company_id")) updateParts.Add("company_id=VALUES(company_id)");
|
|
|
+ updateParts.Add("update_time=VALUES(create_time)");
|
|
|
|
|
|
var sql = $"""
|
|
|
INSERT INTO `{table}`
|
|
|
({string.Join(", ", insertCols.Select(c => $"`{c}`"))})
|
|
|
VALUES
|
|
|
- ({string.Join(", ", insertVals)})
|
|
|
+ {string.Join(",\n ", valueGroups)}
|
|
|
ON DUPLICATE KEY UPDATE
|
|
|
{string.Join(",\n ", updateParts)}
|
|
|
""";
|
|
|
-
|
|
|
- return await _db.Ado.ExecuteCommandAsync(sql, parameters);
|
|
|
+ return (sql, parameters);
|
|
|
}
|
|
|
|
|
|
internal static long? NormalizeFactoryId(long tenantId, long? factoryId) =>
|