MdpDbPullExecutor.cs 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483
  1. using System.Data;
  2. using System.Globalization;
  3. using System.Text.Json;
  4. using System.Text.Json.Serialization;
  5. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  6. using SqlSugar;
  7. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  8. /// <summary>
  9. /// 方式甲:从 mdp_source DB 源按 mdp_entity.source_table_name 增量抽数 → target_table_name(贴源)。
  10. /// </summary>
  11. public sealed class MdpDbPullExecutor : IMdpSourcePullExecutor, ITransient
  12. {
  13. public string SupportedType => "DB_SYNC";
  14. /// <summary>
  15. /// 贴源 raw_data 的时间统一按 MySQL 字面量格式输出;各域转换层用
  16. /// STR_TO_DATE(..., '%Y-%m-%d %H:%i:%s.%f') 解析,ISO 8601 的 T 分隔符会解析失败。
  17. /// </summary>
  18. private static readonly JsonSerializerOptions RawDataJsonOptions = new()
  19. {
  20. Converters =
  21. {
  22. new MysqlLiteralDateTimeConverter(),
  23. new MysqlLiteralDateTimeOffsetConverter()
  24. }
  25. };
  26. private readonly MdpSourceScopeFactory _scopeFactory;
  27. private readonly ISqlSugarClient _db;
  28. private readonly MdpStagingWriter _writer;
  29. public MdpDbPullExecutor(MdpSourceScopeFactory scopeFactory, ISqlSugarClient db, MdpStagingWriter writer)
  30. {
  31. _scopeFactory = scopeFactory;
  32. _db = db;
  33. _writer = writer;
  34. }
  35. public async Task<MdpPullResult> PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default)
  36. {
  37. if (string.IsNullOrWhiteSpace(entity.SourceTableName))
  38. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 source_table_name");
  39. if (string.IsNullOrWhiteSpace(entity.TargetTableName))
  40. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 target_table_name");
  41. var startedAt = DateTime.Now;
  42. var scope = await _scopeFactory.GetScopeAsync(source.SourceCode, cancellationToken);
  43. var batchSize = entity.BatchSize > 0 ? entity.BatchSize : 1000;
  44. var isSqlServer = string.Equals(source.DbType, "SQLSERVER", StringComparison.OrdinalIgnoreCase)
  45. || string.Equals(source.DbType, "MSSQL", StringComparison.OrdinalIgnoreCase);
  46. var useKeyset = ctx.UseKeysetCursor
  47. && !string.IsNullOrWhiteSpace(ctx.CursorColumn)
  48. && !string.IsNullOrWhiteSpace(ctx.TieBreakerColumn);
  49. string? tenantCol = null;
  50. string? factoryCol = null;
  51. if (ctx.TenantId > 0)
  52. {
  53. var cols = await TryGetSourceColumnsAsync(scope, entity.SourceTableName!, isSqlServer, cancellationToken);
  54. tenantCol = FindColumn(cols, "tenant_id", "TenantId", "TenantID");
  55. if (ctx.FactoryId > 0)
  56. factoryCol = FindColumn(cols, "factory_id", "FactoryId");
  57. }
  58. var sql = useKeyset
  59. ? BuildKeysetSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol)
  60. : BuildSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol);
  61. var parameters = new List<SugarParameter>();
  62. if (tenantCol != null)
  63. parameters.Add(new SugarParameter("@scopeTenantId", ctx.TenantId));
  64. if (factoryCol != null)
  65. parameters.Add(new SugarParameter("@scopeFactoryId", ctx.FactoryId));
  66. if (useKeyset)
  67. {
  68. AddKeysetParameters(parameters, ctx);
  69. }
  70. else
  71. {
  72. if (ctx.WindowFrom.HasValue)
  73. parameters.Add(new SugarParameter("@windowFrom", ctx.WindowFrom.Value));
  74. if (!ctx.FullRefresh
  75. && !string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase)
  76. && !string.IsNullOrWhiteSpace(entity.IncrColumn)
  77. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  78. && !string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase))
  79. parameters.Add(new SugarParameter("@cursor", entity.LastCursor));
  80. // ROLLING:用窗口下界;若已有水位且水位晚于下界,则进一步用水位收窄
  81. if (!ctx.FullRefresh
  82. && string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase)
  83. && !string.IsNullOrWhiteSpace(entity.IncrColumn)
  84. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  85. && ctx.WindowFrom.HasValue
  86. && DateTime.TryParse(entity.LastCursor, out var cursorDt)
  87. && cursorDt > ctx.WindowFrom.Value)
  88. parameters.Add(new SugarParameter("@cursor", entity.LastCursor));
  89. }
  90. var table = await scope.Ado.GetDataTableAsync(sql, parameters);
  91. cancellationToken.ThrowIfCancellationRequested();
  92. var written = 0;
  93. string? maxCursor = entity.LastCursor;
  94. var now = DateTime.Now;
  95. string? lastKeysetCursor = null;
  96. string? lastKeysetTie = null;
  97. foreach (DataRow row in table.Rows)
  98. {
  99. cancellationToken.ThrowIfCancellationRequested();
  100. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  101. foreach (DataColumn col in table.Columns)
  102. dict[col.ColumnName] = row[col] == DBNull.Value ? null : row[col];
  103. var sourceRowId = ResolveSourceRowId(dict);
  104. var rawJson = JsonSerializer.Serialize(dict, RawDataJsonOptions);
  105. if (useKeyset)
  106. {
  107. lastKeysetCursor = FormatCursorValue(dict, ctx.CursorColumn!);
  108. lastKeysetTie = FormatCursorValue(dict, ctx.TieBreakerColumn!);
  109. if (!string.IsNullOrEmpty(lastKeysetTie))
  110. maxCursor = EncodeKeysetCursor(lastKeysetCursor, lastKeysetTie);
  111. }
  112. else if (!string.IsNullOrWhiteSpace(entity.IncrColumn) && dict.TryGetValue(entity.IncrColumn, out var incrVal) && incrVal != null)
  113. {
  114. var cursor = incrVal is DateTime dt ? dt.ToString("yyyy-MM-dd HH:mm:ss.fff") : incrVal.ToString();
  115. if (!string.IsNullOrEmpty(cursor) && (maxCursor == null || string.CompareOrdinal(cursor, maxCursor) > 0))
  116. maxCursor = cursor;
  117. }
  118. written += await _writer.UpsertAsync(
  119. source, entity, entity.SourceTableName!, dict, rawJson, sourceRowId, ctx);
  120. }
  121. if (useKeyset && !string.IsNullOrEmpty(lastKeysetTie))
  122. {
  123. ctx.CursorValue = lastKeysetCursor;
  124. ctx.TieBreakerValue = lastKeysetTie;
  125. }
  126. if (!ctx.DeferCursorPersist
  127. && !string.IsNullOrEmpty(maxCursor)
  128. && maxCursor != entity.LastCursor)
  129. {
  130. await _db.Updateable<MdpEntity>()
  131. .SetColumns(x => new MdpEntity
  132. {
  133. LastCursor = maxCursor,
  134. LastSyncTo = now,
  135. UpdateTime = now
  136. })
  137. .Where(x => x.Id == entity.Id)
  138. .ExecuteCommandAsync(cancellationToken);
  139. }
  140. await WriteSyncLogAsync(source, entity, ctx, startedAt, table.Rows.Count, written, null);
  141. var windowHint = useKeyset
  142. ? (ctx.NullTimePhase ? "KEYSET_NULL" : "KEYSET")
  143. : (ctx.SyncWindowType ?? (ctx.FullRefresh ? "FULL" : "INCR"));
  144. return new MdpPullResult
  145. {
  146. RowsPulled = table.Rows.Count,
  147. RowsWritten = written,
  148. NewCursor = maxCursor,
  149. Message = $"OK window={windowHint}" + (ctx.WindowFrom.HasValue ? $" from={ctx.WindowFrom:yyyy-MM-dd HH:mm:ss}" : "")
  150. + (ctx.BootstrapFrom.HasValue ? $" bootstrapFrom={ctx.BootstrapFrom:yyyy-MM-dd HH:mm:ss}" : "")
  151. };
  152. }
  153. private static string BuildSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
  154. {
  155. var table = entity.SourceTableName!.Trim();
  156. // 仅允许简单标识符/schema.table,防注入
  157. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_\.\[\]]+$"))
  158. throw new InvalidOperationException($"非法 source_table_name:{table}");
  159. var predicates = new List<string>();
  160. string orderBy;
  161. if (!string.IsNullOrWhiteSpace(entity.IncrColumn))
  162. {
  163. var incr = entity.IncrColumn.Trim();
  164. if (!System.Text.RegularExpressions.Regex.IsMatch(incr, @"^[A-Za-z0-9_]+$"))
  165. throw new InvalidOperationException($"非法 incr_column:{incr}");
  166. var incrExpr = isSqlServer ? incr : $"`{incr}`";
  167. orderBy = incrExpr;
  168. var isRolling = string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase);
  169. if (!ctx.FullRefresh && isRolling && ctx.WindowFrom.HasValue)
  170. predicates.Add($"{incrExpr} >= @windowFrom");
  171. var useCursor = !ctx.FullRefresh
  172. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  173. && (!isRolling
  174. || (ctx.WindowFrom.HasValue
  175. && DateTime.TryParse(entity.LastCursor, out var cursorDt)
  176. && cursorDt > ctx.WindowFrom.Value));
  177. if (useCursor)
  178. predicates.Add($"{incrExpr} > @cursor");
  179. }
  180. else
  181. {
  182. orderBy = isSqlServer ? "(SELECT NULL)" : "1";
  183. }
  184. AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol);
  185. var where = predicates.Count > 0 ? " WHERE " + string.Join(" AND ", predicates) : "";
  186. var offset = ctx.Offset > 0 ? ctx.Offset : 0;
  187. if (isSqlServer)
  188. {
  189. // SQL Server:OFFSET 需 ORDER BY;无 incr 时用稳定键兜底
  190. if (string.IsNullOrWhiteSpace(entity.IncrColumn))
  191. orderBy = "(SELECT NULL)";
  192. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET {offset} ROWS FETCH NEXT {batchSize} ROWS ONLY";
  193. }
  194. return offset > 0
  195. ? $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize} OFFSET {offset}"
  196. : $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}";
  197. }
  198. private static string BuildKeysetSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
  199. {
  200. var table = entity.SourceTableName!.Trim();
  201. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_\.\[\]]+$"))
  202. throw new InvalidOperationException($"非法 source_table_name:{table}");
  203. var cursorCol = RequireIdent(ctx.CursorColumn, "CursorColumn");
  204. var tieCol = RequireIdent(ctx.TieBreakerColumn, "TieBreakerColumn");
  205. var cursorExpr = isSqlServer ? cursorCol : $"`{cursorCol}`";
  206. var tieExpr = isSqlServer ? tieCol : $"`{tieCol}`";
  207. var predicates = new List<string>();
  208. if (ctx.NullTimePhase)
  209. {
  210. predicates.Add($"{cursorExpr} IS NULL");
  211. if (!string.IsNullOrWhiteSpace(ctx.TieBreakerValue))
  212. predicates.Add($"{tieExpr} > @tieBreaker");
  213. }
  214. else
  215. {
  216. predicates.Add($"{cursorExpr} IS NOT NULL");
  217. if (ctx.BootstrapFrom.HasValue)
  218. predicates.Add($"{cursorExpr} >= @bootstrapFrom");
  219. // keyset lower bound
  220. predicates.Add($"""
  221. (
  222. @cursorTime IS NULL
  223. OR {cursorExpr} > @cursorTime
  224. OR ({cursorExpr} = @cursorTime AND {tieExpr} > @tieBreaker)
  225. )
  226. """);
  227. if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue))
  228. {
  229. predicates.Add($"""
  230. (
  231. {cursorExpr} < @upperTime
  232. OR ({cursorExpr} = @upperTime AND {tieExpr} <= @upperTie)
  233. )
  234. """);
  235. }
  236. }
  237. AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol);
  238. var where = " WHERE " + string.Join(" AND ", predicates);
  239. var orderBy = ctx.NullTimePhase ? tieExpr : $"{cursorExpr}, {tieExpr}";
  240. if (isSqlServer)
  241. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET 0 ROWS FETCH NEXT {batchSize} ROWS ONLY";
  242. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}";
  243. }
  244. private static void AddKeysetParameters(List<SugarParameter> parameters, MdpPullContext ctx)
  245. {
  246. if (ctx.BootstrapFrom.HasValue)
  247. parameters.Add(new SugarParameter("@bootstrapFrom", ctx.BootstrapFrom.Value));
  248. if (ctx.NullTimePhase)
  249. {
  250. parameters.Add(new SugarParameter("@tieBreaker",
  251. string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue));
  252. return;
  253. }
  254. object cursorTime = string.IsNullOrWhiteSpace(ctx.CursorValue)
  255. ? DBNull.Value
  256. : (object)ctx.CursorValue!;
  257. parameters.Add(new SugarParameter("@cursorTime", cursorTime));
  258. parameters.Add(new SugarParameter("@tieBreaker",
  259. string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue));
  260. if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue))
  261. {
  262. parameters.Add(new SugarParameter("@upperTime", ctx.UpperCursorValue));
  263. parameters.Add(new SugarParameter("@upperTie",
  264. string.IsNullOrWhiteSpace(ctx.UpperTieBreakerValue) ? "0" : ctx.UpperTieBreakerValue));
  265. }
  266. }
  267. private static void AppendScopePredicates(List<string> predicates, bool isSqlServer, string? tenantCol, string? factoryCol)
  268. {
  269. if (!string.IsNullOrWhiteSpace(tenantCol))
  270. predicates.Add($"{QuoteIdent(tenantCol, isSqlServer)} = @scopeTenantId");
  271. if (!string.IsNullOrWhiteSpace(factoryCol))
  272. {
  273. var f = QuoteIdent(factoryCol, isSqlServer);
  274. predicates.Add($"COALESCE(NULLIF({f}, 0), 1) = @scopeFactoryId");
  275. }
  276. }
  277. private static string QuoteIdent(string name, bool isSqlServer) => isSqlServer ? name : $"`{name}`";
  278. private static string? FindColumn(IReadOnlyCollection<string> columns, params string[] names)
  279. {
  280. foreach (var n in names)
  281. {
  282. var hit = columns.FirstOrDefault(c => string.Equals(c, n, StringComparison.OrdinalIgnoreCase));
  283. if (hit != null) return hit;
  284. }
  285. return null;
  286. }
  287. private static async Task<IReadOnlyCollection<string>> TryGetSourceColumnsAsync(
  288. ISqlSugarClient scope, string table, bool isSqlServer, CancellationToken cancellationToken)
  289. {
  290. try
  291. {
  292. var simple = table.Contains('.') ? table.Split('.').Last().Trim('[', ']') : table.Trim('[', ']');
  293. if (!System.Text.RegularExpressions.Regex.IsMatch(simple, @"^[A-Za-z0-9_]+$"))
  294. return Array.Empty<string>();
  295. if (isSqlServer)
  296. {
  297. var rows = await scope.Ado.SqlQueryAsync<string>(
  298. "SELECT name FROM sys.columns WHERE object_id = OBJECT_ID(@t)",
  299. new SugarParameter("@t", simple));
  300. return rows ?? new List<string>();
  301. }
  302. var mysql = await scope.Ado.SqlQueryAsync<string>(
  303. """
  304. SELECT COLUMN_NAME FROM information_schema.COLUMNS
  305. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  306. """,
  307. new SugarParameter("@t", simple));
  308. return mysql ?? new List<string>();
  309. }
  310. catch
  311. {
  312. return Array.Empty<string>();
  313. }
  314. }
  315. private static string RequireIdent(string? name, string label)
  316. {
  317. var v = (name ?? "").Trim();
  318. if (!System.Text.RegularExpressions.Regex.IsMatch(v, @"^[A-Za-z0-9_]+$"))
  319. throw new InvalidOperationException($"非法 {label}:{name}");
  320. return v;
  321. }
  322. private static string? FormatCursorValue(Dictionary<string, object?> dict, string column)
  323. {
  324. if (!dict.TryGetValue(column, out var raw) || raw == null || raw is DBNull)
  325. return null;
  326. return raw switch
  327. {
  328. DateTime dt => dt.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture),
  329. DateTimeOffset dto => dto.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture),
  330. _ => raw.ToString()
  331. };
  332. }
  333. internal static string EncodeKeysetCursor(string? cursor, string tieBreaker)
  334. => System.Text.Json.JsonSerializer.Serialize(new KeysetCursorDto
  335. {
  336. Cursor = cursor,
  337. TieBreaker = tieBreaker
  338. });
  339. internal static bool TryDecodeKeysetCursor(string? raw, out string? cursor, out string? tieBreaker)
  340. {
  341. cursor = null;
  342. tieBreaker = null;
  343. if (string.IsNullOrWhiteSpace(raw))
  344. return false;
  345. try
  346. {
  347. var dto = System.Text.Json.JsonSerializer.Deserialize<KeysetCursorDto>(raw);
  348. if (dto == null || string.IsNullOrWhiteSpace(dto.TieBreaker))
  349. return false;
  350. cursor = dto.Cursor;
  351. tieBreaker = dto.TieBreaker;
  352. return true;
  353. }
  354. catch
  355. {
  356. return false;
  357. }
  358. }
  359. private sealed class KeysetCursorDto
  360. {
  361. public string? Cursor { get; set; }
  362. public string? TieBreaker { get; set; }
  363. }
  364. private static string ResolveSourceRowId(Dictionary<string, object?> dict)
  365. {
  366. // InvTransHist 等表常有空字符串 ID,若优先取 ID 会导致 source_row_id="" 唯一键互相覆盖。
  367. // 跳过 null/空白后,再落到 RecID 等真实键。
  368. foreach (var key in new[]
  369. {
  370. "RecID", "RecId", "recid", "Recid",
  371. "id", "Id", "ID",
  372. "noid", "billno", "BillNo", "djbh", "Djbh"
  373. })
  374. {
  375. if (!dict.TryGetValue(key, out var v) || v is null || v is DBNull)
  376. continue;
  377. var s = v.ToString();
  378. if (!string.IsNullOrWhiteSpace(s))
  379. return s;
  380. }
  381. return Guid.NewGuid().ToString("N");
  382. }
  383. private async Task WriteSyncLogAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, DateTime startedAt, int pulled, int written, string? error)
  384. {
  385. try
  386. {
  387. var endedAt = DateTime.Now;
  388. var syncType = ctx.FullRefresh || string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase)
  389. ? "FULL"
  390. : "INCR";
  391. await _db.Ado.ExecuteCommandAsync(@"
  392. INSERT INTO mdp_sync_log
  393. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type,
  394. sync_start, sync_end, duration_ms, rows_read, rows_written, status, error_message)
  395. VALUES
  396. (@tenantId, @entityId, @sourceCode, @entityName, @batchId, @syncType, 'AUTO',
  397. @syncStart, @syncEnd, @durationMs, @rowsRead, @rowsWritten, @status, @error)",
  398. new SugarParameter("@tenantId", ctx.TenantId),
  399. new SugarParameter("@entityId", entity.Id),
  400. new SugarParameter("@sourceCode", source.SourceCode),
  401. new SugarParameter("@entityName", string.IsNullOrWhiteSpace(entity.EntityName) ? entity.EntityCode : entity.EntityName),
  402. new SugarParameter("@batchId", ctx.BatchId),
  403. new SugarParameter("@syncType", syncType),
  404. new SugarParameter("@syncStart", startedAt),
  405. new SugarParameter("@syncEnd", endedAt),
  406. new SugarParameter("@durationMs", (int)(endedAt - startedAt).TotalMilliseconds),
  407. new SugarParameter("@rowsRead", pulled),
  408. new SugarParameter("@rowsWritten", written),
  409. new SugarParameter("@status", string.IsNullOrEmpty(error) ? "SUCCESS" : "FAILED"),
  410. new SugarParameter("@error", error));
  411. }
  412. catch
  413. {
  414. // 日志表结构可能与种子不一致时不阻断主流程
  415. }
  416. }
  417. private sealed class MysqlLiteralDateTimeConverter : JsonConverter<DateTime>
  418. {
  419. public override DateTime Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
  420. => reader.GetDateTime();
  421. public override void Write(Utf8JsonWriter writer, DateTime value, JsonSerializerOptions options)
  422. => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture));
  423. }
  424. private sealed class MysqlLiteralDateTimeOffsetConverter : JsonConverter<DateTimeOffset>
  425. {
  426. public override DateTimeOffset Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
  427. => reader.GetDateTimeOffset();
  428. public override void Write(Utf8JsonWriter writer, DateTimeOffset value, JsonSerializerOptions options)
  429. => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture));
  430. }
  431. }