- Add per-cycle correlation ID ([cyc:xxxxxxxx]) threaded through Python logs and SQL audit tables for end-to-end traceability of any divergence. - New ProductionDataBaseSync.SyncApplyRunLog table + @CycleID on usp_SyncApply for per-table apply auditing (pending/merged/deleted/applied, dead/error counts, duration). Audit write isolated in its own TRY/CATCH outside txn. - Capture/cleanup phases now emit full detail: Insert->Delete downgrade warning with record_id/log_id, dedup_skipped as apply-stall signal, per-file summary, and access log-id ranges on cleanup. - Queue health check surfaces error/dead rows with recent samples instead of silent accumulation (previously the top cause of data divergence). - sql_writer uses INSERT...SELECT...WHERE NOT EXISTS for observable dedup; idle cycles lowered to DEBUG with periodic heartbeat. - Backward compatible: old proc callers still work (CycleID nullable; legacy coarse-grained logging with one-time notice). Excluded from this commit: CODE.md, build_code_doc.py (doc generation).
204 lines
10 KiB
Transact-SQL
204 lines
10 KiB
Transact-SQL
-- usp_SyncApply: set-based apply of pending SyncQueue rows.
|
|
-- Per distinct (TargetSchema, TargetTable) it builds column lists from sys.columns
|
|
-- (excluding the ID key, computed, identity and rowversion columns) and runs a
|
|
-- dynamic-SQL MERGE (last-write-wins by SourceLogID DESC) for Insert/Update and a
|
|
-- DELETE for the last op = Delete. SET IDENTITY_INSERT ON preserves Access PKs.
|
|
--
|
|
-- CORRECTNESS NOTE: a SINGLE ranked CTE (rn=1 per RecordID, over ALL pending ops
|
|
-- regardless of OperateType) is read by BOTH the upsert and delete branches.
|
|
-- Earlier this was two independent ranked CTEs (one over Insert/Update, one over
|
|
-- Delete), which let a stale Delete outrank a newer Insert for the same RecordID
|
|
-- and silently drop a row whose true last op was an Insert. Routing off the one
|
|
-- winning row's OperateType guarantees only the genuine last op wins.
|
|
--
|
|
-- AUDIT NOTE: when ProductionDataBaseSync.SyncApplyRunLog exists (see
|
|
-- sql/04_sync_apply_runlog.sql), one audit row is written per target table and
|
|
-- per invocation: pending/distinct counts going in, the MERGE/DELETE rowcounts,
|
|
-- how many queue rows were flipped to applied (or to error/dead on failure),
|
|
-- the error message, the duration, and the caller-supplied @CycleID that ties
|
|
-- the row to the Python service log's [cyc:...] lines. The audit insert is
|
|
-- best-effort (wrapped in its own TRY/CATCH, outside the data transaction) and
|
|
-- can never fail the apply itself. @CycleID defaults to NULL so the legacy
|
|
-- single-parameter EXEC keeps working.
|
|
|
|
CREATE OR ALTER PROCEDURE ProductionDataBaseSync.usp_SyncApply
|
|
@MaxRetries INT = 5,
|
|
@CycleID NVARCHAR(40) = NULL
|
|
AS
|
|
BEGIN
|
|
SET NOCOUNT ON;
|
|
|
|
DECLARE @sch NVARCHAR(128), @tbl NVARCHAR(255), @FullName NVARCHAR(514);
|
|
DECLARE @sql NVARCHAR(MAX), @cols NVARCHAR(MAX), @upd NVARCHAR(MAX), @ins NVARCHAR(MAX);
|
|
DECLARE @hasRunLog bit =
|
|
CASE WHEN OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog', 'U') IS NOT NULL
|
|
THEN 1 ELSE 0 END;
|
|
DECLARE @t0 DATETIME2, @pending INT, @distinctRecs INT,
|
|
@merged INT, @deleted INT, @applied INT,
|
|
@errCnt INT, @deadCnt INT,
|
|
@errMsg NVARCHAR(MAX), @outcome VARCHAR(10);
|
|
|
|
-- Re-queue rows still under the retry budget.
|
|
UPDATE ProductionDataBaseSync.SyncQueue
|
|
SET Status = 'pending'
|
|
WHERE Status = 'error' AND RetryCount < @MaxRetries;
|
|
|
|
DECLARE cur CURSOR LOCAL FAST_FORWARD FOR
|
|
SELECT DISTINCT TargetSchema, TargetTable
|
|
FROM ProductionDataBaseSync.SyncQueue
|
|
WHERE Status = 'pending';
|
|
|
|
OPEN cur;
|
|
FETCH NEXT FROM cur INTO @sch, @tbl;
|
|
|
|
WHILE @@FETCH_STATUS = 0
|
|
BEGIN
|
|
SET @FullName = QUOTENAME(@sch) + N'.' + QUOTENAME(@tbl);
|
|
|
|
-- Per-table audit state. Explicit reset every iteration: local
|
|
-- variables keep their previous value across cursor loops.
|
|
SELECT @t0 = SYSDATETIME(),
|
|
@merged = NULL, @deleted = NULL, @applied = NULL,
|
|
@errCnt = NULL, @deadCnt = NULL,
|
|
@errMsg = NULL, @outcome = 'ok',
|
|
@cols = NULL, @upd = NULL, @ins = NULL;
|
|
|
|
SELECT @pending = COUNT(*),
|
|
@distinctRecs = COUNT(DISTINCT RecordID)
|
|
FROM ProductionDataBaseSync.SyncQueue
|
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
|
|
|
BEGIN TRY
|
|
BEGIN TRAN;
|
|
-- @cols: comma-quoted column names (for INSERT target list)
|
|
-- @upd: tgt.col = JSON_VALUE(src.RowData,'$.col')
|
|
-- @ins: JSON_VALUE(src.RowData,'$.col') (for INSERT VALUES)
|
|
-- NOTE: every JSON path key is wrapped in double quotes ('$."col"').
|
|
-- This is mandatory for non-ASCII column names (e.g. Chinese names
|
|
-- like 名字) and harmless for ASCII names, so we quote unconditionally.
|
|
SELECT
|
|
@cols = STRING_AGG(QUOTENAME(c.name), N',')
|
|
WITHIN GROUP (ORDER BY c.column_id),
|
|
@upd = STRING_AGG(QUOTENAME(c.name) + N'=JSON_VALUE(src.RowData,''$."' + c.name + N'"'')', N',')
|
|
WITHIN GROUP (ORDER BY c.column_id),
|
|
@ins = STRING_AGG(N'JSON_VALUE(src.RowData,''$."' + c.name + N'"'')', N',')
|
|
WITHIN GROUP (ORDER BY c.column_id)
|
|
FROM sys.columns c
|
|
WHERE c.object_id = OBJECT_ID(@FullName)
|
|
AND c.is_computed = 0
|
|
AND c.is_identity = 0
|
|
AND TYPE_NAME(c.system_type_id) <> 'timestamp'
|
|
AND c.name <> 'ID';
|
|
|
|
IF @cols IS NOT NULL
|
|
BEGIN
|
|
-- Upsert: winners (rn=1) whose winning OperateType is
|
|
-- Insert/Update and carry RowData. The ranked CTE covers ALL
|
|
-- pending ops so the rn=1 row is the true last op for that
|
|
-- RecordID; an earlier Delete can no longer outrank a newer
|
|
-- Insert for the same RecordID.
|
|
-- AUDIT: @@ROWCOUNT is captured into @MergedOut IMMEDIATELY
|
|
-- after the MERGE -- SET IDENTITY_INSERT (like any SET option)
|
|
-- resets @@ROWCOUNT, so the order below is load-bearing.
|
|
SET @sql = N'SET IDENTITY_INSERT ' + @FullName + N' ON;'
|
|
+ N';WITH ranked AS ('
|
|
+ N' SELECT RecordID, OperateType, RowData,'
|
|
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
|
+ N' FROM ProductionDataBaseSync.SyncQueue'
|
|
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
|
+ N')'
|
|
+ N'MERGE ' + @FullName + N' WITH (HOLDLOCK) AS tgt'
|
|
+ N' USING ('
|
|
+ N' SELECT RecordID, RowData FROM ranked'
|
|
+ N' WHERE rn=1 AND OperateType IN (''Insert'',''Update'') AND RowData IS NOT NULL'
|
|
+ N') AS src ON tgt.ID = TRY_CAST(src.RecordID AS int)'
|
|
+ N' WHEN MATCHED THEN UPDATE SET ' + @upd
|
|
+ N' WHEN NOT MATCHED THEN INSERT (ID,' + @cols + N')'
|
|
+ N' VALUES (TRY_CAST(src.RecordID AS int),' + @ins + N');'
|
|
+ N'SET @MergedOut = @@ROWCOUNT;'
|
|
+ N'SET IDENTITY_INSERT ' + @FullName + N' OFF;';
|
|
EXEC sp_executesql @sql,
|
|
N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@MergedOut INT OUTPUT',
|
|
@sch, @tbl, @MergedOut = @merged OUTPUT;
|
|
|
|
-- Delete: winners (rn=1) whose winning OperateType is Delete.
|
|
-- Same ranked CTE — only the genuine last op can be a delete.
|
|
SET @sql = N';WITH ranked AS ('
|
|
+ N' SELECT RecordID, OperateType,'
|
|
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
|
+ N' FROM ProductionDataBaseSync.SyncQueue'
|
|
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
|
+ N')'
|
|
+ N'DELETE t FROM ' + @FullName + N' t'
|
|
+ N' JOIN (SELECT RecordID FROM ranked WHERE rn=1 AND OperateType=''Delete'') d'
|
|
+ N' ON t.ID = TRY_CAST(d.RecordID AS int);'
|
|
+ N'SET @DeletedOut = @@ROWCOUNT;';
|
|
EXEC sp_executesql @sql,
|
|
N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@DeletedOut INT OUTPUT',
|
|
@sch, @tbl, @DeletedOut = @deleted OUTPUT;
|
|
END
|
|
|
|
UPDATE ProductionDataBaseSync.SyncQueue
|
|
SET Status = 'applied', AppliedAt = SYSDATETIME()
|
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
|
SET @applied = @@ROWCOUNT;
|
|
|
|
COMMIT;
|
|
END TRY
|
|
BEGIN CATCH
|
|
-- ROLLBACK unwinds the per-table transaction (and, if the caller
|
|
-- started one, the caller's too). A proc that errors under an outer
|
|
-- transaction is expected to leave the txn uncommittable; the caller
|
|
-- issues its own rollback. We roll back here so the partial per-table
|
|
-- work is discarded before marking rows as 'error'/'dead'.
|
|
IF @@TRANCOUNT > 0 ROLLBACK;
|
|
SELECT @errMsg = ERROR_MESSAGE(), @outcome = 'error';
|
|
|
|
-- Split of the previous single CASE update so the audit row can
|
|
-- report exactly how many rows died vs. how many will be retried.
|
|
-- Net effect on SyncQueue is identical.
|
|
UPDATE ProductionDataBaseSync.SyncQueue
|
|
SET Status = 'dead', RetryCount = RetryCount + 1, ErrorMsg = @errMsg
|
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending'
|
|
AND RetryCount + 1 >= @MaxRetries;
|
|
SET @deadCnt = @@ROWCOUNT;
|
|
|
|
UPDATE ProductionDataBaseSync.SyncQueue
|
|
SET Status = 'error', RetryCount = RetryCount + 1, ErrorMsg = @errMsg
|
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
|
SET @errCnt = @@ROWCOUNT;
|
|
END CATCH
|
|
|
|
-- Best-effort audit row: sits outside the data transaction (after
|
|
-- COMMIT/ROLLBACK) so it survives either outcome, and its own failure
|
|
-- can never break the apply loop.
|
|
IF @hasRunLog = 1
|
|
BEGIN
|
|
BEGIN TRY
|
|
INSERT ProductionDataBaseSync.SyncApplyRunLog
|
|
(CycleID, TargetSchema, TargetTable,
|
|
PendingCount, DistinctRecords,
|
|
MergedCount, DeletedCount, AppliedCount,
|
|
ErrorCount, DeadCount,
|
|
Outcome, ErrorMsg, StartedAt, DurationMs)
|
|
VALUES
|
|
(@CycleID, @sch, @tbl,
|
|
@pending, @distinctRecs,
|
|
@merged, @deleted, @applied,
|
|
@errCnt, @deadCnt,
|
|
@outcome, @errMsg, @t0,
|
|
DATEDIFF(millisecond, @t0, SYSDATETIME()));
|
|
END TRY
|
|
BEGIN CATCH
|
|
PRINT 'SyncApplyRunLog insert failed: ' + ERROR_MESSAGE();
|
|
END CATCH
|
|
END
|
|
|
|
FETCH NEXT FROM cur INTO @sch, @tbl;
|
|
END
|
|
|
|
CLOSE cur;
|
|
DEALLOCATE cur;
|
|
END
|
|
GO
|