Compare commits
5 Commits
7d11eddc7a
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2ce03d21d9 | ||
|
|
75fb6a3a01 | ||
|
|
d71b7eab62 | ||
|
|
4179d3e232 | ||
|
|
d8a423c983 |
17
main.py
17
main.py
@@ -17,6 +17,7 @@ was launched or whether the venv already has ``src`` on its path.
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import argparse
|
import argparse
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
|
|
||||||
@@ -27,9 +28,10 @@ from sync.config import load_config
|
|||||||
from sync.logging_setup import setup_logging
|
from sync.logging_setup import setup_logging
|
||||||
from sync.fullsync import full_sync
|
from sync.fullsync import full_sync
|
||||||
from sync import service
|
from sync import service
|
||||||
from sync.compare import compare, any_mismatch, format_report
|
from sync.compare import compare, any_mismatch, format_report, write_report
|
||||||
|
|
||||||
CONFIG_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.yaml")
|
CONFIG_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.yaml")
|
||||||
|
log = logging.getLogger("main")
|
||||||
|
|
||||||
|
|
||||||
def _parse_args(argv):
|
def _parse_args(argv):
|
||||||
@@ -99,11 +101,14 @@ def main(argv=None) -> int:
|
|||||||
if args.command == "compare":
|
if args.command == "compare":
|
||||||
results = compare(cfg, granularity=args.granularity,
|
results = compare(cfg, granularity=args.granularity,
|
||||||
db_filter=args.db, table_filter=args.table)
|
db_filter=args.db, table_filter=args.table)
|
||||||
report = format_report(results, args.granularity)
|
# Persist the report as a dated log file under the logging directory
|
||||||
print(report)
|
# (logs/ by default); --report still allows an extra custom path.
|
||||||
if args.report:
|
log_dir = os.path.dirname((cfg.logging or {}).get("path", "sync.log")) or "."
|
||||||
with open(args.report, "w", encoding="utf-8") as f:
|
written = write_report(results, args.granularity, log_dir=log_dir,
|
||||||
f.write(report + "\n")
|
extra_path=args.report)
|
||||||
|
print(format_report(results, args.granularity))
|
||||||
|
log.info("compare finished (granularity=%s) report=%s mismatch=%s",
|
||||||
|
args.granularity, written, any_mismatch(results))
|
||||||
return 1 if any_mismatch(results) else 0
|
return 1 if any_mismatch(results) else 0
|
||||||
|
|
||||||
return 2 # unreachable: argparse requires a subcommand
|
return 2 # unreachable: argparse requires a subcommand
|
||||||
|
|||||||
3
run_compare_ids.cmd
Normal file
3
run_compare_ids.cmd
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
@echo off
|
||||||
|
cd /d C:\Users\peng\Projects\ProductionDataBaseSync_DataMacro
|
||||||
|
.venv\Scripts\python.exe main.py compare --granularity ids >> logs\compare_task_stdout.log 2>&1
|
||||||
12
sql/00_schema.sql
Normal file
12
sql/00_schema.sql
Normal file
@@ -0,0 +1,12 @@
|
|||||||
|
-- ProductionDataBaseSync: dedicated schema that groups every object owned by
|
||||||
|
-- the Access -> SQL Server sync (SyncQueue staging, SyncLogArchive audit store,
|
||||||
|
-- and the usp_SyncApply procedure). Keeping them out of dbo makes ownership and
|
||||||
|
-- housekeeping explicit.
|
||||||
|
-- Idempotent: CREATE SCHEMA must be the only statement in its batch, so it is
|
||||||
|
-- wrapped in EXEC() behind an existence check.
|
||||||
|
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.schemas WHERE name = 'ProductionDataBaseSync')
|
||||||
|
BEGIN
|
||||||
|
EXEC('CREATE SCHEMA ProductionDataBaseSync');
|
||||||
|
END
|
||||||
|
GO
|
||||||
@@ -1,9 +1,10 @@
|
|||||||
-- SyncQueue: staging table for the Access -> SQL Server one-way sync.
|
-- SyncQueue: staging table for the Access -> SQL Server one-way sync.
|
||||||
|
-- Lives under the ProductionDataBaseSync schema (run 00_schema.sql first).
|
||||||
-- Idempotent: safe to re-run (table created only if absent; indexes only if absent).
|
-- Idempotent: safe to re-run (table created only if absent; indexes only if absent).
|
||||||
|
|
||||||
IF OBJECT_ID('dbo.SyncQueue', 'U') IS NULL
|
IF OBJECT_ID('ProductionDataBaseSync.SyncQueue', 'U') IS NULL
|
||||||
BEGIN
|
BEGIN
|
||||||
CREATE TABLE dbo.SyncQueue (
|
CREATE TABLE ProductionDataBaseSync.SyncQueue (
|
||||||
QueueID bigint IDENTITY(1,1) NOT NULL,
|
QueueID bigint IDENTITY(1,1) NOT NULL,
|
||||||
SourceFile nvarchar(255) NOT NULL,
|
SourceFile nvarchar(255) NOT NULL,
|
||||||
SourceTable nvarchar(255) NOT NULL,
|
SourceTable nvarchar(255) NOT NULL,
|
||||||
@@ -18,6 +19,7 @@ BEGIN
|
|||||||
ErrorMsg nvarchar(max) NULL,
|
ErrorMsg nvarchar(max) NULL,
|
||||||
CapturedAt datetime2 NOT NULL CONSTRAINT DF_SyncQueue_Captured DEFAULT sysdatetime(),
|
CapturedAt datetime2 NOT NULL CONSTRAINT DF_SyncQueue_Captured DEFAULT sysdatetime(),
|
||||||
AppliedAt datetime2 NULL,
|
AppliedAt datetime2 NULL,
|
||||||
|
CleanedAt datetime2 NULL,
|
||||||
CONSTRAINT PK_SyncQueue PRIMARY KEY CLUSTERED (QueueID)
|
CONSTRAINT PK_SyncQueue PRIMARY KEY CLUSTERED (QueueID)
|
||||||
);
|
);
|
||||||
END
|
END
|
||||||
@@ -25,38 +27,27 @@ GO
|
|||||||
|
|
||||||
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
WHERE name = 'UX_SyncQueue_Dedup'
|
WHERE name = 'UX_SyncQueue_Dedup'
|
||||||
AND object_id = OBJECT_ID('dbo.SyncQueue'))
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
|
||||||
BEGIN
|
BEGIN
|
||||||
CREATE UNIQUE INDEX UX_SyncQueue_Dedup
|
CREATE UNIQUE INDEX UX_SyncQueue_Dedup
|
||||||
ON dbo.SyncQueue(SourceFile, SourceTable, SourceLogID);
|
ON ProductionDataBaseSync.SyncQueue(SourceFile, SourceTable, SourceLogID);
|
||||||
END
|
END
|
||||||
GO
|
GO
|
||||||
|
|
||||||
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
WHERE name = 'IX_SyncQueue_Pending'
|
WHERE name = 'IX_SyncQueue_Pending'
|
||||||
AND object_id = OBJECT_ID('dbo.SyncQueue'))
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
|
||||||
BEGIN
|
BEGIN
|
||||||
CREATE INDEX IX_SyncQueue_Pending
|
CREATE INDEX IX_SyncQueue_Pending
|
||||||
ON dbo.SyncQueue(Status, TargetSchema, TargetTable);
|
ON ProductionDataBaseSync.SyncQueue(Status, TargetSchema, TargetTable);
|
||||||
END
|
|
||||||
GO
|
|
||||||
|
|
||||||
-- Cleanup bookkeeping: track when a queue row's Access log counterpart has
|
|
||||||
-- been physically deleted, so cleanup never re-deletes the same IDs and the
|
|
||||||
-- table can be purged to bound its growth.
|
|
||||||
IF NOT EXISTS (SELECT 1 FROM sys.columns
|
|
||||||
WHERE object_id = OBJECT_ID('dbo.SyncQueue')
|
|
||||||
AND name = 'CleanedAt')
|
|
||||||
BEGIN
|
|
||||||
ALTER TABLE dbo.SyncQueue ADD CleanedAt datetime2 NULL;
|
|
||||||
END
|
END
|
||||||
GO
|
GO
|
||||||
|
|
||||||
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
WHERE name = 'IX_SyncQueue_Cleaned'
|
WHERE name = 'IX_SyncQueue_Cleaned'
|
||||||
AND object_id = OBJECT_ID('dbo.SyncQueue'))
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
|
||||||
BEGIN
|
BEGIN
|
||||||
CREATE INDEX IX_SyncQueue_Cleaned
|
CREATE INDEX IX_SyncQueue_Cleaned
|
||||||
ON dbo.SyncQueue(Status, CleanedAt);
|
ON ProductionDataBaseSync.SyncQueue(Status, CleanedAt);
|
||||||
END
|
END
|
||||||
GO
|
GO
|
||||||
|
|||||||
@@ -10,24 +10,42 @@
|
|||||||
-- Delete), which let a stale Delete outrank a newer Insert for the same RecordID
|
-- 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
|
-- 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.
|
-- 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 dbo.usp_SyncApply
|
CREATE OR ALTER PROCEDURE ProductionDataBaseSync.usp_SyncApply
|
||||||
@MaxRetries INT = 5
|
@MaxRetries INT = 5,
|
||||||
|
@CycleID NVARCHAR(40) = NULL
|
||||||
AS
|
AS
|
||||||
BEGIN
|
BEGIN
|
||||||
SET NOCOUNT ON;
|
SET NOCOUNT ON;
|
||||||
|
|
||||||
DECLARE @sch NVARCHAR(128), @tbl NVARCHAR(255), @FullName NVARCHAR(514);
|
DECLARE @sch NVARCHAR(128), @tbl NVARCHAR(255), @FullName NVARCHAR(514);
|
||||||
DECLARE @sql NVARCHAR(MAX), @cols NVARCHAR(MAX), @upd NVARCHAR(MAX), @ins NVARCHAR(MAX);
|
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.
|
-- Re-queue rows still under the retry budget.
|
||||||
UPDATE dbo.SyncQueue
|
UPDATE ProductionDataBaseSync.SyncQueue
|
||||||
SET Status = 'pending'
|
SET Status = 'pending'
|
||||||
WHERE Status = 'error' AND RetryCount < @MaxRetries;
|
WHERE Status = 'error' AND RetryCount < @MaxRetries;
|
||||||
|
|
||||||
DECLARE cur CURSOR LOCAL FAST_FORWARD FOR
|
DECLARE cur CURSOR LOCAL FAST_FORWARD FOR
|
||||||
SELECT DISTINCT TargetSchema, TargetTable
|
SELECT DISTINCT TargetSchema, TargetTable
|
||||||
FROM dbo.SyncQueue
|
FROM ProductionDataBaseSync.SyncQueue
|
||||||
WHERE Status = 'pending';
|
WHERE Status = 'pending';
|
||||||
|
|
||||||
OPEN cur;
|
OPEN cur;
|
||||||
@@ -37,6 +55,19 @@ BEGIN
|
|||||||
BEGIN
|
BEGIN
|
||||||
SET @FullName = QUOTENAME(@sch) + N'.' + QUOTENAME(@tbl);
|
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 TRY
|
||||||
BEGIN TRAN;
|
BEGIN TRAN;
|
||||||
-- @cols: comma-quoted column names (for INSERT target list)
|
-- @cols: comma-quoted column names (for INSERT target list)
|
||||||
@@ -66,11 +97,14 @@ BEGIN
|
|||||||
-- pending ops so the rn=1 row is the true last op for that
|
-- pending ops so the rn=1 row is the true last op for that
|
||||||
-- RecordID; an earlier Delete can no longer outrank a newer
|
-- RecordID; an earlier Delete can no longer outrank a newer
|
||||||
-- Insert for the same RecordID.
|
-- 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;'
|
SET @sql = N'SET IDENTITY_INSERT ' + @FullName + N' ON;'
|
||||||
+ N';WITH ranked AS ('
|
+ N';WITH ranked AS ('
|
||||||
+ N' SELECT RecordID, OperateType, RowData,'
|
+ N' SELECT RecordID, OperateType, RowData,'
|
||||||
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
||||||
+ N' FROM dbo.SyncQueue'
|
+ N' FROM ProductionDataBaseSync.SyncQueue'
|
||||||
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
||||||
+ N')'
|
+ N')'
|
||||||
+ N'MERGE ' + @FullName + N' WITH (HOLDLOCK) AS tgt'
|
+ N'MERGE ' + @FullName + N' WITH (HOLDLOCK) AS tgt'
|
||||||
@@ -81,28 +115,33 @@ BEGIN
|
|||||||
+ N' WHEN MATCHED THEN UPDATE SET ' + @upd
|
+ N' WHEN MATCHED THEN UPDATE SET ' + @upd
|
||||||
+ N' WHEN NOT MATCHED THEN INSERT (ID,' + @cols + N')'
|
+ N' WHEN NOT MATCHED THEN INSERT (ID,' + @cols + N')'
|
||||||
+ N' VALUES (TRY_CAST(src.RecordID AS int),' + @ins + N');'
|
+ N' VALUES (TRY_CAST(src.RecordID AS int),' + @ins + N');'
|
||||||
|
+ N'SET @MergedOut = @@ROWCOUNT;'
|
||||||
+ N'SET IDENTITY_INSERT ' + @FullName + N' OFF;';
|
+ N'SET IDENTITY_INSERT ' + @FullName + N' OFF;';
|
||||||
EXEC sp_executesql @sql,
|
EXEC sp_executesql @sql,
|
||||||
N'@sch NVARCHAR(128),@tbl NVARCHAR(255)', @sch, @tbl;
|
N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@MergedOut INT OUTPUT',
|
||||||
|
@sch, @tbl, @MergedOut = @merged OUTPUT;
|
||||||
|
|
||||||
-- Delete: winners (rn=1) whose winning OperateType is Delete.
|
-- Delete: winners (rn=1) whose winning OperateType is Delete.
|
||||||
-- Same ranked CTE — only the genuine last op can be a delete.
|
-- Same ranked CTE — only the genuine last op can be a delete.
|
||||||
SET @sql = N';WITH ranked AS ('
|
SET @sql = N';WITH ranked AS ('
|
||||||
+ N' SELECT RecordID, OperateType,'
|
+ N' SELECT RecordID, OperateType,'
|
||||||
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
|
||||||
+ N' FROM dbo.SyncQueue'
|
+ N' FROM ProductionDataBaseSync.SyncQueue'
|
||||||
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
|
||||||
+ N')'
|
+ N')'
|
||||||
+ N'DELETE t FROM ' + @FullName + N' t'
|
+ N'DELETE t FROM ' + @FullName + N' t'
|
||||||
+ N' JOIN (SELECT RecordID FROM ranked WHERE rn=1 AND OperateType=''Delete'') d'
|
+ N' JOIN (SELECT RecordID FROM ranked WHERE rn=1 AND OperateType=''Delete'') d'
|
||||||
+ N' ON t.ID = TRY_CAST(d.RecordID AS int);';
|
+ N' ON t.ID = TRY_CAST(d.RecordID AS int);'
|
||||||
|
+ N'SET @DeletedOut = @@ROWCOUNT;';
|
||||||
EXEC sp_executesql @sql,
|
EXEC sp_executesql @sql,
|
||||||
N'@sch NVARCHAR(128),@tbl NVARCHAR(255)', @sch, @tbl;
|
N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@DeletedOut INT OUTPUT',
|
||||||
|
@sch, @tbl, @DeletedOut = @deleted OUTPUT;
|
||||||
END
|
END
|
||||||
|
|
||||||
UPDATE dbo.SyncQueue
|
UPDATE ProductionDataBaseSync.SyncQueue
|
||||||
SET Status = 'applied', AppliedAt = SYSDATETIME()
|
SET Status = 'applied', AppliedAt = SYSDATETIME()
|
||||||
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
||||||
|
SET @applied = @@ROWCOUNT;
|
||||||
|
|
||||||
COMMIT;
|
COMMIT;
|
||||||
END TRY
|
END TRY
|
||||||
@@ -113,13 +152,48 @@ BEGIN
|
|||||||
-- issues its own rollback. We roll back here so the partial per-table
|
-- issues its own rollback. We roll back here so the partial per-table
|
||||||
-- work is discarded before marking rows as 'error'/'dead'.
|
-- work is discarded before marking rows as 'error'/'dead'.
|
||||||
IF @@TRANCOUNT > 0 ROLLBACK;
|
IF @@TRANCOUNT > 0 ROLLBACK;
|
||||||
UPDATE dbo.SyncQueue
|
SELECT @errMsg = ERROR_MESSAGE(), @outcome = 'error';
|
||||||
SET Status = CASE WHEN RetryCount + 1 >= @MaxRetries THEN 'dead' ELSE 'error' END,
|
|
||||||
RetryCount = RetryCount + 1,
|
-- Split of the previous single CASE update so the audit row can
|
||||||
ErrorMsg = ERROR_MESSAGE()
|
-- 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';
|
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
|
||||||
|
SET @errCnt = @@ROWCOUNT;
|
||||||
END CATCH
|
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;
|
FETCH NEXT FROM cur INTO @sch, @tbl;
|
||||||
END
|
END
|
||||||
|
|
||||||
|
|||||||
71
sql/03_sync_log_archive.sql
Normal file
71
sql/03_sync_log_archive.sql
Normal file
@@ -0,0 +1,71 @@
|
|||||||
|
-- SyncLogArchive: permanent, append-only audit store for every Access change-log
|
||||||
|
-- row consumed by the incremental sync. This is the durable evidence layer that
|
||||||
|
-- SyncQueue is NOT: SyncQueue is a transient work queue (purged 24h after a row
|
||||||
|
-- is cleaned) and it only stores the *processed* OperateType, so a downgrade
|
||||||
|
-- (e.g. capture turning an Insert into a Delete when the row is momentarily
|
||||||
|
-- unreadable) erases the original intent. This table preserves both the ORIGINAL
|
||||||
|
-- operate type recorded by the Access data macro and the PROCESSED type actually
|
||||||
|
-- sent to SQL Server, plus the original Access log timestamp and the row payload.
|
||||||
|
--
|
||||||
|
-- With this in place, cases like the 14287 incident (a real Insert applied to SQL
|
||||||
|
-- as a Delete) stay fully reconstructible: OriginalOperateType != ProcessedOperateType
|
||||||
|
-- flags exactly where the pipeline diverged from the source.
|
||||||
|
--
|
||||||
|
-- Retention: PERMANENT. No purge job touches this table. SyncQueue keeps its
|
||||||
|
-- short-lived queue role; this table keeps history. Lives under the
|
||||||
|
-- ProductionDataBaseSync schema (run 00_schema.sql first).
|
||||||
|
-- Idempotent: safe to re-run.
|
||||||
|
|
||||||
|
IF OBJECT_ID('ProductionDataBaseSync.SyncLogArchive', 'U') IS NULL
|
||||||
|
BEGIN
|
||||||
|
CREATE TABLE ProductionDataBaseSync.SyncLogArchive (
|
||||||
|
ArchiveID bigint IDENTITY(1,1) NOT NULL,
|
||||||
|
SourceFile nvarchar(255) NOT NULL,
|
||||||
|
SourceTable nvarchar(255) NOT NULL,
|
||||||
|
SourceLogID bigint NOT NULL,
|
||||||
|
RecordID nvarchar(50) NOT NULL,
|
||||||
|
TargetSchema nvarchar(128) NOT NULL,
|
||||||
|
TargetTable nvarchar(255) NOT NULL,
|
||||||
|
OriginalOperateType varchar(10) NOT NULL, -- as recorded by the Access data macro
|
||||||
|
ProcessedOperateType varchar(10) NOT NULL, -- as actually sent to SyncQueue / SQL
|
||||||
|
RowData nvarchar(max) NULL, -- captured row payload (NULL for Delete)
|
||||||
|
OriginalTime datetime2 NULL, -- Access TableChangeLog.Time (previously discarded)
|
||||||
|
CapturedAt datetime2 NOT NULL
|
||||||
|
CONSTRAINT DF_SyncLogArchive_Captured DEFAULT sysdatetime(),
|
||||||
|
CONSTRAINT PK_SyncLogArchive PRIMARY KEY CLUSTERED (ArchiveID)
|
||||||
|
);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- Dedup: one archive row per source log entry. Capture may re-run the same log
|
||||||
|
-- row if a prior cycle's apply failed (the Access log is only deleted after a
|
||||||
|
-- successful apply), so the write path uses IF NOT EXISTS on these keys.
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'UX_SyncLogArchive_Dedup'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
|
||||||
|
BEGIN
|
||||||
|
CREATE UNIQUE INDEX UX_SyncLogArchive_Dedup
|
||||||
|
ON ProductionDataBaseSync.SyncLogArchive(SourceFile, SourceTable, SourceLogID);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- Evidence lookup by table + record (e.g. "show every log ever seen for ID 14287").
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'IX_SyncLogArchive_Record'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
|
||||||
|
BEGIN
|
||||||
|
CREATE INDEX IX_SyncLogArchive_Record
|
||||||
|
ON ProductionDataBaseSync.SyncLogArchive(SourceTable, RecordID);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- Fast filter for the anomaly the archive exists to catch: original != processed.
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'IX_SyncLogArchive_Divergence'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
|
||||||
|
BEGIN
|
||||||
|
CREATE INDEX IX_SyncLogArchive_Divergence
|
||||||
|
ON ProductionDataBaseSync.SyncLogArchive(OriginalOperateType, ProcessedOperateType)
|
||||||
|
INCLUDE (SourceFile, SourceTable, RecordID, CapturedAt);
|
||||||
|
END
|
||||||
|
GO
|
||||||
75
sql/04_sync_apply_runlog.sql
Normal file
75
sql/04_sync_apply_runlog.sql
Normal file
@@ -0,0 +1,75 @@
|
|||||||
|
-- SyncApplyRunLog: per-invocation, per-target-table audit of what usp_SyncApply
|
||||||
|
-- actually did. This closes the biggest observability gap in the pipeline: the
|
||||||
|
-- apply phase used to be a black box ("apply done") -- queue rows could flip
|
||||||
|
-- to 'error' or 'dead' with no trace in the service log, and there was no
|
||||||
|
-- record of how many rows a MERGE/DELETE touched at any point in time. With
|
||||||
|
-- this table, "what did apply do to table X around time T, and did it fail?"
|
||||||
|
-- is a single indexed query, and CycleID joins each row back to the exact
|
||||||
|
-- [cyc:xxxxxxxx] lines in the Python service log.
|
||||||
|
--
|
||||||
|
-- Written by usp_SyncApply (sql/02_sync_apply.sql) as a best-effort insert per
|
||||||
|
-- (invocation, target table). Idle cycles write nothing (the proc's cursor
|
||||||
|
-- only visits tables that have pending rows), so growth tracks real change
|
||||||
|
-- traffic, not poll frequency. Retention: unmanaged by default; if it ever
|
||||||
|
-- grows large, purge by StartedAt, e.g.
|
||||||
|
-- DELETE FROM ProductionDataBaseSync.SyncApplyRunLog
|
||||||
|
-- WHERE StartedAt < DATEADD(day, -90, SYSDATETIME());
|
||||||
|
--
|
||||||
|
-- Lives under the ProductionDataBaseSync schema (run 00_schema.sql first).
|
||||||
|
-- Idempotent: safe to re-run. Deploy alongside the updated 02_sync_apply.sql;
|
||||||
|
-- ordering is forgiving either way (the proc checks for this table and skips
|
||||||
|
-- the audit insert when it is absent).
|
||||||
|
|
||||||
|
IF OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog', 'U') IS NULL
|
||||||
|
BEGIN
|
||||||
|
CREATE TABLE ProductionDataBaseSync.SyncApplyRunLog (
|
||||||
|
RunLogID bigint IDENTITY(1,1) NOT NULL,
|
||||||
|
CycleID nvarchar(40) NULL, -- correlation id from the Python service ([cyc:...])
|
||||||
|
TargetSchema nvarchar(128) NOT NULL,
|
||||||
|
TargetTable nvarchar(255) NOT NULL,
|
||||||
|
PendingCount int NOT NULL, -- pending queue rows seen for this table
|
||||||
|
DistinctRecords int NULL, -- distinct RecordIDs among them
|
||||||
|
MergedCount int NULL, -- rows affected by the MERGE (insert + update)
|
||||||
|
DeletedCount int NULL, -- rows affected by the DELETE
|
||||||
|
AppliedCount int NULL, -- queue rows flipped to 'applied'
|
||||||
|
ErrorCount int NULL, -- queue rows flipped to 'error' (will retry)
|
||||||
|
DeadCount int NULL, -- queue rows flipped to 'dead' (retries exhausted)
|
||||||
|
Outcome varchar(10) NOT NULL, -- 'ok' | 'error'
|
||||||
|
ErrorMsg nvarchar(max) NULL,
|
||||||
|
StartedAt datetime2 NOT NULL
|
||||||
|
CONSTRAINT DF_SyncApplyRunLog_Started DEFAULT sysdatetime(),
|
||||||
|
DurationMs int NULL,
|
||||||
|
CONSTRAINT PK_SyncApplyRunLog PRIMARY KEY CLUSTERED (RunLogID)
|
||||||
|
);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- Join back to the Python service log of one cycle.
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'IX_SyncApplyRunLog_Cycle'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
|
||||||
|
BEGIN
|
||||||
|
CREATE INDEX IX_SyncApplyRunLog_Cycle
|
||||||
|
ON ProductionDataBaseSync.SyncApplyRunLog(CycleID);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- "What happened to this table around time T?"
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'IX_SyncApplyRunLog_Table'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
|
||||||
|
BEGIN
|
||||||
|
CREATE INDEX IX_SyncApplyRunLog_Table
|
||||||
|
ON ProductionDataBaseSync.SyncApplyRunLog(TargetSchema, TargetTable, StartedAt);
|
||||||
|
END
|
||||||
|
GO
|
||||||
|
|
||||||
|
-- Fast scan for failed applies.
|
||||||
|
IF NOT EXISTS (SELECT 1 FROM sys.indexes
|
||||||
|
WHERE name = 'IX_SyncApplyRunLog_Outcome'
|
||||||
|
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
|
||||||
|
BEGIN
|
||||||
|
CREATE INDEX IX_SyncApplyRunLog_Outcome
|
||||||
|
ON ProductionDataBaseSync.SyncApplyRunLog(Outcome, StartedAt);
|
||||||
|
END
|
||||||
|
GO
|
||||||
@@ -94,7 +94,9 @@ class AccessReader:
|
|||||||
``cursor.rowcount``), so the caller can report honest counts. A no-op
|
``cursor.rowcount``), so the caller can report honest counts. A no-op
|
||||||
(returns 0) when ``ids`` is empty. Retries with linear backoff because
|
(returns 0) when ``ids`` is empty. Retries with linear backoff because
|
||||||
the live client may briefly hold a page lock on ``TableChangeLog``.
|
the live client may briefly hold a page lock on ``TableChangeLog``.
|
||||||
Raises on final lock failure.
|
Each retry is logged at WARNING (previously a silent sleep) and the
|
||||||
|
final lock failure at ERROR before raising, so lock churn on the
|
||||||
|
production files is visible in the audit trail.
|
||||||
"""
|
"""
|
||||||
if not ids:
|
if not ids:
|
||||||
return 0
|
return 0
|
||||||
@@ -121,8 +123,20 @@ class AccessReader:
|
|||||||
msg = str(e)
|
msg = str(e)
|
||||||
is_lock = "被锁定" in msg or "-1102" in msg
|
is_lock = "被锁定" in msg or "-1102" in msg
|
||||||
if is_lock and attempt < retries - 1:
|
if is_lock and attempt < retries - 1:
|
||||||
|
log.warning(
|
||||||
|
"TableChangeLog delete lock contention "
|
||||||
|
"(attempt %d/%d) db=%s ids=%s..%s -- retrying: %s",
|
||||||
|
attempt + 1, retries, self.db_path,
|
||||||
|
chunk[0], chunk[-1], msg[:200],
|
||||||
|
)
|
||||||
time.sleep(0.2 * (attempt + 1))
|
time.sleep(0.2 * (attempt + 1))
|
||||||
else:
|
else:
|
||||||
|
if is_lock:
|
||||||
|
log.error(
|
||||||
|
"TableChangeLog delete still locked after %d "
|
||||||
|
"attempts db=%s ids=%s..%s -- giving up",
|
||||||
|
retries, self.db_path, chunk[0], chunk[-1],
|
||||||
|
)
|
||||||
raise
|
raise
|
||||||
return total
|
return total
|
||||||
|
|
||||||
|
|||||||
@@ -1,39 +1,175 @@
|
|||||||
|
"""Capture phase: read Access change-log rows and stage them into SyncQueue.
|
||||||
|
|
||||||
|
Every consumed ``TableChangeLog`` row is appended to the permanent audit store
|
||||||
|
(``SyncLogArchive``) BEFORE it is enqueued, so the original evidence survives
|
||||||
|
cleanup. On top of that, this module now emits a detailed audit trail to the
|
||||||
|
service log:
|
||||||
|
|
||||||
|
- WARNING for every operate-type downgrade (Insert/Update -> Delete because
|
||||||
|
the source row was unreadable at capture time), with record id, log id and
|
||||||
|
the Access log timestamp -- the exact event class behind the 14287 incident,
|
||||||
|
previously invisible in the text log;
|
||||||
|
- WARNING for unknown operate types, with enough identity (log id / record id
|
||||||
|
/ time) to locate and repair the offending log row manually;
|
||||||
|
- a per-file INFO summary: rows read, newly enqueued, dedup-skipped
|
||||||
|
(re-capture after a failed apply -- a symptom worth noticing), downgraded,
|
||||||
|
out-of-scope, unknown ops, the processed log-ID range and a per-operation
|
||||||
|
breakdown;
|
||||||
|
- DEBUG detail for individual out-of-scope and dedup skips.
|
||||||
|
"""
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
import json, logging
|
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
|
||||||
from .access_reader import AccessReader
|
from .access_reader import AccessReader
|
||||||
from .sql_writer import SqlWriter, QueueRow
|
from .sql_writer import SqlWriter, QueueRow, ArchiveRow
|
||||||
from .config import FileMapping, SyncConfig
|
from .config import FileMapping, SyncConfig
|
||||||
from .targets import is_synced_table
|
from .targets import is_synced_table
|
||||||
|
|
||||||
log = logging.getLogger(__name__)
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
def capture_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter, cfg: SyncConfig) -> int:
|
|
||||||
|
@dataclass
|
||||||
|
class CaptureStats:
|
||||||
|
"""Counters for one capture pass (per file, or aggregated per cycle)."""
|
||||||
|
|
||||||
|
read: int = 0 # change-log rows read from Access
|
||||||
|
enqueued: int = 0 # rows newly inserted into SyncQueue
|
||||||
|
dedup_skipped: int = 0 # already queued (re-capture after a failed apply)
|
||||||
|
downgraded: int = 0 # Insert/Update downgraded to Delete
|
||||||
|
out_of_scope: int = 0 # log rows for tables outside the sync scope
|
||||||
|
unknown_op: int = 0 # log rows with an unrecognised OperateType
|
||||||
|
min_log_id: int | None = None
|
||||||
|
max_log_id: int | None = None
|
||||||
|
ops: dict[str, int] = field(default_factory=dict) # processed op -> count
|
||||||
|
|
||||||
|
def note_op(self, op: str) -> None:
|
||||||
|
self.ops[op] = self.ops.get(op, 0) + 1
|
||||||
|
|
||||||
|
def merge(self, other: "CaptureStats") -> None:
|
||||||
|
"""Fold *other* into this instance (cycle-level aggregation)."""
|
||||||
|
self.read += other.read
|
||||||
|
self.enqueued += other.enqueued
|
||||||
|
self.dedup_skipped += other.dedup_skipped
|
||||||
|
self.downgraded += other.downgraded
|
||||||
|
self.out_of_scope += other.out_of_scope
|
||||||
|
self.unknown_op += other.unknown_op
|
||||||
|
for k, v in other.ops.items():
|
||||||
|
self.ops[k] = self.ops.get(k, 0) + v
|
||||||
|
for attr, pick in (("min_log_id", min), ("max_log_id", max)):
|
||||||
|
a, b = getattr(self, attr), getattr(other, attr)
|
||||||
|
if a is None:
|
||||||
|
setattr(self, attr, b)
|
||||||
|
elif b is not None:
|
||||||
|
setattr(self, attr, pick(a, b))
|
||||||
|
|
||||||
|
def ops_str(self) -> str:
|
||||||
|
return ",".join(f"{k}={v}" for k, v in sorted(self.ops.items())) or "-"
|
||||||
|
|
||||||
|
|
||||||
|
def capture_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter,
|
||||||
|
cfg: SyncConfig) -> CaptureStats:
|
||||||
|
"""Capture one file's pending change-log rows. Returns detailed stats.
|
||||||
|
|
||||||
|
NOTE: previously returned a bare int (rows enqueued); that count is now
|
||||||
|
``stats.enqueued``. ``sync.service.cycle`` is the only in-repo caller and
|
||||||
|
has been updated accordingly.
|
||||||
|
"""
|
||||||
rows = reader.read_log(cfg.runtime.capture_batch_size)
|
rows = reader.read_log(cfg.runtime.capture_batch_size)
|
||||||
n = 0
|
st = CaptureStats(read=len(rows))
|
||||||
|
if rows: # read_log orders by ID ascending
|
||||||
|
st.min_log_id, st.max_log_id = rows[0].id, rows[-1].id
|
||||||
|
|
||||||
for lr in rows:
|
for lr in rows:
|
||||||
if not is_synced_table(fm, lr.table_name):
|
if not is_synced_table(fm, lr.table_name):
|
||||||
|
st.out_of_scope += 1
|
||||||
|
log.debug("capture skip (out of scope) file=%s table=%s log_id=%s",
|
||||||
|
fm.file, lr.table_name, lr.id)
|
||||||
continue
|
continue
|
||||||
op = lr.operate_type
|
op = lr.operate_type
|
||||||
row_data = None
|
row_data = None
|
||||||
if op in ("Insert", "Update"):
|
if op in ("Insert", "Update"):
|
||||||
d = reader.read_row(lr.table_name, lr.record_id)
|
d = reader.read_row(lr.table_name, lr.record_id)
|
||||||
if d is None:
|
if d is None:
|
||||||
op = "Delete" # 行已删,降级
|
# 行已删(或此刻不可读),降级为 Delete。这是数据差异排查的
|
||||||
|
# 头号嫌疑事件(参见 14287 事件),必须在文本日志显式留痕,
|
||||||
|
# 而不只是写入 SyncLogArchive。
|
||||||
|
op = "Delete"
|
||||||
|
st.downgraded += 1
|
||||||
|
log.warning(
|
||||||
|
"capture DOWNGRADE %s->Delete file=%s table=%s "
|
||||||
|
"record_id=%s log_id=%s log_time=%s "
|
||||||
|
"(source row unreadable at capture time; original intent "
|
||||||
|
"preserved in SyncLogArchive.OriginalOperateType)",
|
||||||
|
lr.operate_type, fm.file, lr.table_name,
|
||||||
|
lr.record_id, lr.id, lr.time,
|
||||||
|
)
|
||||||
else:
|
else:
|
||||||
row_data = json.dumps(d, ensure_ascii=False)
|
row_data = json.dumps(d, ensure_ascii=False)
|
||||||
elif op != "Delete":
|
elif op != "Delete":
|
||||||
log.warning("unknown OperateType %r in %s log %s", op, fm.file, lr.id)
|
st.unknown_op += 1
|
||||||
|
log.warning(
|
||||||
|
"capture UNKNOWN OperateType %r file=%s table=%s "
|
||||||
|
"record_id=%s log_id=%s log_time=%s -- row skipped; it will "
|
||||||
|
"be re-read every cycle until removed from TableChangeLog",
|
||||||
|
op, fm.file, lr.table_name, lr.record_id, lr.id, lr.time,
|
||||||
|
)
|
||||||
continue
|
continue
|
||||||
|
target_schema = fm.schema
|
||||||
|
target_table = fm.target_table(lr.table_name)
|
||||||
|
# Persist the ORIGINAL log entry to the permanent audit store BEFORE the
|
||||||
|
# queue insert (and long before cleanup deletes the Access log). This
|
||||||
|
# keeps both the source operate type (lr.operate_type) and the processed
|
||||||
|
# one (op) so a downgrade like Insert->Delete stays reconstructible.
|
||||||
|
writer.insert_archive_row(ArchiveRow(
|
||||||
|
source_file=fm.file,
|
||||||
|
source_table=lr.table_name,
|
||||||
|
source_log_id=lr.id,
|
||||||
|
record_id=lr.record_id,
|
||||||
|
target_schema=target_schema,
|
||||||
|
target_table=target_table,
|
||||||
|
original_operate_type=lr.operate_type,
|
||||||
|
processed_operate_type=op,
|
||||||
|
row_data=row_data,
|
||||||
|
original_time=lr.time,
|
||||||
|
))
|
||||||
qr = QueueRow(
|
qr = QueueRow(
|
||||||
source_file=fm.file,
|
source_file=fm.file,
|
||||||
source_table=lr.table_name,
|
source_table=lr.table_name,
|
||||||
record_id=lr.record_id,
|
record_id=lr.record_id,
|
||||||
target_schema=fm.schema,
|
target_schema=target_schema,
|
||||||
target_table=fm.target_table(lr.table_name),
|
target_table=target_table,
|
||||||
source_log_id=lr.id,
|
source_log_id=lr.id,
|
||||||
operate_type=op,
|
operate_type=op,
|
||||||
row_data=row_data,
|
row_data=row_data,
|
||||||
)
|
)
|
||||||
writer.insert_queue_row(qr)
|
if writer.insert_queue_row(qr):
|
||||||
n += 1
|
st.enqueued += 1
|
||||||
return n
|
st.note_op(op)
|
||||||
|
else:
|
||||||
|
# Dedup hit: this log row was already staged by an earlier cycle
|
||||||
|
# whose apply failed (the Access log row is only deleted after a
|
||||||
|
# successful apply). A persistently non-zero dedup count therefore
|
||||||
|
# points straight at a stuck apply -- see the queue-health WARNINGs
|
||||||
|
# emitted by sync.service.cycle.
|
||||||
|
st.dedup_skipped += 1
|
||||||
|
log.debug(
|
||||||
|
"capture dedup-skip (already queued) file=%s table=%s "
|
||||||
|
"log_id=%s record_id=%s op=%s",
|
||||||
|
fm.file, lr.table_name, lr.id, lr.record_id, op,
|
||||||
|
)
|
||||||
|
|
||||||
|
if st.read:
|
||||||
|
log.info(
|
||||||
|
"capture file=%s read=%d enqueued=%d dedup_skipped=%d "
|
||||||
|
"downgraded=%d out_of_scope=%d unknown_op=%d "
|
||||||
|
"log_ids=%s..%s ops={%s}",
|
||||||
|
fm.file, st.read, st.enqueued, st.dedup_skipped, st.downgraded,
|
||||||
|
st.out_of_scope, st.unknown_op, st.min_log_id, st.max_log_id,
|
||||||
|
st.ops_str(),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
log.debug("capture file=%s: change log empty", fm.file)
|
||||||
|
return st
|
||||||
|
|||||||
@@ -1,10 +1,16 @@
|
|||||||
"""Cleanup phase: delete applied Access log rows.
|
"""Cleanup phase: delete applied Access log rows.
|
||||||
|
|
||||||
After ``dbo.usp_SyncApply`` flips queue rows to ``applied``, those rows'
|
After ``usp_SyncApply`` flips queue rows to ``applied``, those rows'
|
||||||
``SourceLogID`` values are no longer needed on the Access side. This module
|
``SourceLogID`` values are no longer needed on the Access side. This module
|
||||||
asks the writer which log IDs have been applied for a given source file and
|
asks the writer which log IDs have been applied for a given source file and
|
||||||
deletes them from ``TableChangeLog`` via the reader, in batches with lock
|
deletes them from ``TableChangeLog`` via the reader, in batches with lock
|
||||||
retry. The delete is the only mutation that touches the Access side.
|
retry. The delete is the only mutation that touches the Access side.
|
||||||
|
|
||||||
|
Audit trail: the INFO line records the exact log-ID range removed from each
|
||||||
|
file, and a WARNING is raised when fewer rows were deleted than expected --
|
||||||
|
i.e. some applied log IDs were already absent from the Access log (removed
|
||||||
|
externally, or by a previously interrupted run), which is worth knowing when
|
||||||
|
reconstructing what happened around a divergence.
|
||||||
"""
|
"""
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
import logging
|
import logging
|
||||||
@@ -28,9 +34,25 @@ def cleanup_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter, cfg:
|
|||||||
"""
|
"""
|
||||||
ids = writer.applied_log_ids(fm.file)
|
ids = writer.applied_log_ids(fm.file)
|
||||||
if not ids:
|
if not ids:
|
||||||
|
log.debug("cleanup file=%s: no applied rows to clean", fm.file)
|
||||||
return 0
|
return 0
|
||||||
|
log.debug("cleanup file=%s: deleting %d applied log rows ids=%s..%s",
|
||||||
|
fm.file, len(ids), ids[0], ids[-1])
|
||||||
deleted = reader.delete_log_ids(
|
deleted = reader.delete_log_ids(
|
||||||
ids, cfg.runtime.cleanup_batch_size, cfg.runtime.cleanup_lock_retries
|
ids, cfg.runtime.cleanup_batch_size, cfg.runtime.cleanup_lock_retries
|
||||||
)
|
)
|
||||||
|
if deleted != len(ids):
|
||||||
|
log.warning(
|
||||||
|
"cleanup file=%s: deleted %d of %d applied log rows "
|
||||||
|
"(some log IDs were already absent from Access -- removed "
|
||||||
|
"externally or by an earlier interrupted run)",
|
||||||
|
fm.file, deleted, len(ids),
|
||||||
|
)
|
||||||
writer.mark_cleaned(fm.file, ids)
|
writer.mark_cleaned(fm.file, ids)
|
||||||
|
if deleted:
|
||||||
|
log.info(
|
||||||
|
"cleanup file=%s: removed %d access log rows ids=%s..%s; "
|
||||||
|
"queue rows marked cleaned",
|
||||||
|
fm.file, deleted, ids[0], ids[-1],
|
||||||
|
)
|
||||||
return deleted
|
return deleted
|
||||||
|
|||||||
@@ -11,7 +11,9 @@ as full sync -- empirically confirming the two stay aligned.
|
|||||||
"""
|
"""
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import datetime as _dt
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field
|
||||||
|
|
||||||
from .config import FileMapping, SyncConfig
|
from .config import FileMapping, SyncConfig
|
||||||
@@ -139,3 +141,52 @@ def format_report(results: list[TableResult], granularity: str = "count") -> str
|
|||||||
else:
|
else:
|
||||||
lines.append(f"{base} access={r.access_count} sql={r.sql_count} [{r.status.upper()}]")
|
lines.append(f"{base} access={r.access_count} sql={r.sql_count} [{r.status.upper()}]")
|
||||||
return "\n".join(lines)
|
return "\n".join(lines)
|
||||||
|
|
||||||
|
|
||||||
|
def summarize(results: list[TableResult]) -> dict:
|
||||||
|
"""Tally the outcome buckets: match / mismatch / skipped / error."""
|
||||||
|
s = {"match": 0, "mismatch": 0, "skipped": 0, "error": 0}
|
||||||
|
for r in results:
|
||||||
|
s[r.status] = s.get(r.status, 0) + 1
|
||||||
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
def write_report(results: list[TableResult], granularity: str = "count",
|
||||||
|
log_dir: str = ".", run_dt: _dt.datetime | None = None,
|
||||||
|
extra_path: str | None = None) -> str:
|
||||||
|
"""Render the full report and persist it as a dated log file under *log_dir*.
|
||||||
|
|
||||||
|
Always writes ``<log_dir>/compare_<granularity>_<YYYY-MM-DD>.log``,
|
||||||
|
overwriting the day's previous run (one report per day; the logging system's
|
||||||
|
per-day archival later moves yesterday's file into ``logs/Archive/``). When
|
||||||
|
*extra_path* is given (the ``--report`` CLI option) the identical content is
|
||||||
|
also written there for backward compatibility. Returns the primary log path.
|
||||||
|
"""
|
||||||
|
run_dt = run_dt or _dt.datetime.now()
|
||||||
|
stats = summarize(results)
|
||||||
|
body = format_report(results, granularity)
|
||||||
|
header = (
|
||||||
|
"============================================================\n"
|
||||||
|
" 数据一致性核对报告 / Compare Report\n"
|
||||||
|
f" 生成时间 : {run_dt.strftime('%Y-%m-%d %H:%M:%S')}\n"
|
||||||
|
f" 粒度 : {granularity}\n"
|
||||||
|
f" 比对表合计 : {stats['match'] + stats['mismatch']}\n"
|
||||||
|
f" 一致 MATCH : {stats['match']}\n"
|
||||||
|
f" 不一致 MISMATCH : {stats['mismatch']}\n"
|
||||||
|
f" 跳过 SKIPPED : {stats['skipped']}\n"
|
||||||
|
f" 错误 ERROR : {stats['error']}\n"
|
||||||
|
"============================================================\n"
|
||||||
|
)
|
||||||
|
report = header + body + "\n"
|
||||||
|
|
||||||
|
log_dir = log_dir or "."
|
||||||
|
os.makedirs(log_dir, exist_ok=True)
|
||||||
|
primary = os.path.join(
|
||||||
|
log_dir, f"compare_{granularity}_{run_dt.strftime('%Y-%m-%d')}.log"
|
||||||
|
)
|
||||||
|
with open(primary, "w", encoding="utf-8") as f:
|
||||||
|
f.write(report)
|
||||||
|
if extra_path:
|
||||||
|
with open(extra_path, "w", encoding="utf-8") as f:
|
||||||
|
f.write(report)
|
||||||
|
return primary
|
||||||
|
|||||||
@@ -5,7 +5,13 @@ from pydantic import BaseModel, ConfigDict, Field
|
|||||||
|
|
||||||
class SqlServerConfig(BaseModel):
|
class SqlServerConfig(BaseModel):
|
||||||
conn_str: str
|
conn_str: str
|
||||||
sync_queue_table: str = "dbo.SyncQueue"
|
sync_queue_table: str = "ProductionDataBaseSync.SyncQueue"
|
||||||
|
archive_table: str = "ProductionDataBaseSync.SyncLogArchive"
|
||||||
|
apply_proc: str = "ProductionDataBaseSync.usp_SyncApply"
|
||||||
|
# Per-invocation, per-table apply audit written by usp_SyncApply
|
||||||
|
# (see sql/04_sync_apply_runlog.sql). Default matches the shipped script;
|
||||||
|
# existing config.yaml files need no change.
|
||||||
|
apply_runlog_table: str = "ProductionDataBaseSync.SyncApplyRunLog"
|
||||||
|
|
||||||
class AccessConfig(BaseModel):
|
class AccessConfig(BaseModel):
|
||||||
model_config = ConfigDict(coerce_numbers_to_str=True)
|
model_config = ConfigDict(coerce_numbers_to_str=True)
|
||||||
@@ -21,6 +27,9 @@ class RuntimeConfig(BaseModel):
|
|||||||
cleanup_batch_size: int = 200
|
cleanup_batch_size: int = 200
|
||||||
cleanup_lock_retries: int = 3
|
cleanup_lock_retries: int = 3
|
||||||
cleaned_retention_hours: int = 24
|
cleaned_retention_hours: int = 24
|
||||||
|
# Idle cycles now log at DEBUG; the service emits an INFO heartbeat at this
|
||||||
|
# interval while idle so a quiet log still proves the service is alive.
|
||||||
|
idle_heartbeat_seconds: int = 600
|
||||||
|
|
||||||
class FileMapping(BaseModel):
|
class FileMapping(BaseModel):
|
||||||
model_config = ConfigDict(coerce_numbers_to_str=True)
|
model_config = ConfigDict(coerce_numbers_to_str=True)
|
||||||
|
|||||||
@@ -2,10 +2,124 @@
|
|||||||
|
|
||||||
Configures the root logger with a RotatingFileHandler (10 MB x 5, UTF-8) plus
|
Configures the root logger with a RotatingFileHandler (10 MB x 5, UTF-8) plus
|
||||||
a console StreamHandler. The level and log path come from ``cfg.logging``.
|
a console StreamHandler. The level and log path come from ``cfg.logging``.
|
||||||
|
|
||||||
|
Every record written to the log file carries a cycle correlation id
|
||||||
|
(``[cyc:xxxxxxxx]``): ``sync.service.cycle`` allocates one per pass via
|
||||||
|
``set_cycle_id`` and ``_CycleIdFilter`` injects it into each record, so every
|
||||||
|
capture/apply/cleanup line of one pass -- and the matching
|
||||||
|
``SyncApplyRunLog.CycleID`` rows on SQL Server -- can be correlated with a
|
||||||
|
single grep. Outside a cycle (fullsync, compare, startup) the field is ``-``.
|
||||||
|
|
||||||
|
Log files are managed per-day: at startup any previously produced log
|
||||||
|
(including the project's own ``sync.log`` and NSSM's ``nssm_*.log`` captures)
|
||||||
|
is relocated into an ``Archive/`` subfolder next to the active log. Where a log
|
||||||
|
lacks a timestamp, a ``-YYYY-MM-DD`` suffix is added (derived from its first
|
||||||
|
log line, falling back to mtime) so historical files carry a date. The log
|
||||||
|
root therefore only ever shows the current day's ``sync.log``.
|
||||||
"""
|
"""
|
||||||
|
import contextvars
|
||||||
|
import datetime
|
||||||
import logging
|
import logging
|
||||||
import logging.handlers
|
import logging.handlers
|
||||||
import os
|
import os
|
||||||
|
import re
|
||||||
|
import shutil
|
||||||
|
|
||||||
|
|
||||||
|
_LOG_DATE_RE = re.compile(r"^(\d{4}-\d{2}-\d{2})")
|
||||||
|
|
||||||
|
# Current cycle correlation id ("-" when not inside a service cycle).
|
||||||
|
_cycle_id: contextvars.ContextVar[str] = contextvars.ContextVar(
|
||||||
|
"sync_cycle_id", default="-"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def set_cycle_id(cycle_id: str | None) -> None:
|
||||||
|
"""Set (or clear, with ``None``) the id stamped on every log record.
|
||||||
|
|
||||||
|
Called by ``sync.service.cycle`` at the start/end of each pass. The same
|
||||||
|
id is passed to ``usp_SyncApply`` so SQL-side ``SyncApplyRunLog`` rows can
|
||||||
|
be joined back to the exact log lines of the cycle that produced them.
|
||||||
|
"""
|
||||||
|
_cycle_id.set(cycle_id or "-")
|
||||||
|
|
||||||
|
|
||||||
|
class _CycleIdFilter(logging.Filter):
|
||||||
|
"""Inject the current cycle id into every record as ``record.cycle``.
|
||||||
|
|
||||||
|
Attached to the handlers (not the logger) so records emitted through any
|
||||||
|
module logger -- capture, cleanup, access_reader, ... -- are covered.
|
||||||
|
"""
|
||||||
|
|
||||||
|
def filter(self, record: logging.LogRecord) -> bool: # noqa: A003
|
||||||
|
record.cycle = _cycle_id.get()
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
def _first_line_date(path: str) -> str | None:
|
||||||
|
"""Best-effort extraction of the first log line's YYYY-MM-DD date."""
|
||||||
|
try:
|
||||||
|
with open(path, "r", encoding="utf-8", errors="replace") as f:
|
||||||
|
for line in f:
|
||||||
|
m = _LOG_DATE_RE.match(line)
|
||||||
|
if m:
|
||||||
|
return m.group(1)
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _archive_completed_logs(log_dir: str, active_name: str):
|
||||||
|
"""Move already-produced logs out of *log_dir* into *log_dir*/Archive.
|
||||||
|
|
||||||
|
The active log (*active_name*) is left in place only when it belongs to the
|
||||||
|
current day; an older ``sync.log`` is archived (with a ``-YYYY-MM-DD``
|
||||||
|
suffix) so a fresh one can be opened. NSSM's own ``nssm_*.log`` captures are
|
||||||
|
timestamped already and are moved as-is. Moves are best-effort: files locked
|
||||||
|
by another process (e.g. NSSM's live handles) are skipped.
|
||||||
|
"""
|
||||||
|
archive_dir = os.path.join(log_dir, "Archive")
|
||||||
|
os.makedirs(archive_dir, exist_ok=True)
|
||||||
|
today = datetime.date.today()
|
||||||
|
|
||||||
|
for name in os.listdir(log_dir):
|
||||||
|
src = os.path.join(log_dir, name)
|
||||||
|
if not os.path.isfile(src):
|
||||||
|
continue
|
||||||
|
if name == "Archive":
|
||||||
|
continue
|
||||||
|
if not (name.endswith(".log") or ".log." in name):
|
||||||
|
continue
|
||||||
|
# Keep NSSM's live, currently-open handles in place.
|
||||||
|
if name in ("nssm_stderr.log", "nssm_stdout.log"):
|
||||||
|
continue
|
||||||
|
|
||||||
|
if name == active_name:
|
||||||
|
# Only archive the active log if it is from a previous day.
|
||||||
|
log_date = _first_line_date(src)
|
||||||
|
log_date = (
|
||||||
|
datetime.date.fromisoformat(log_date)
|
||||||
|
if log_date
|
||||||
|
else datetime.date.fromtimestamp(os.path.getmtime(src))
|
||||||
|
)
|
||||||
|
if log_date >= today:
|
||||||
|
continue # today's log: keep appending
|
||||||
|
new_name = f"sync-{log_date.strftime('%Y-%m-%d')}.log"
|
||||||
|
else:
|
||||||
|
new_name = name
|
||||||
|
|
||||||
|
dst = os.path.join(archive_dir, new_name)
|
||||||
|
if os.path.exists(dst):
|
||||||
|
stem, ext = os.path.splitext(new_name)
|
||||||
|
i = 2
|
||||||
|
while os.path.exists(os.path.join(archive_dir, f"{stem}({i}){ext}")):
|
||||||
|
i += 1
|
||||||
|
dst = os.path.join(archive_dir, f"{stem}({i}){ext}")
|
||||||
|
try:
|
||||||
|
shutil.move(src, dst)
|
||||||
|
except OSError:
|
||||||
|
# Locked by another process (e.g. NSSM holding the file open).
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def setup_logging(cfg_dict: dict | None):
|
def setup_logging(cfg_dict: dict | None):
|
||||||
@@ -13,18 +127,45 @@ def setup_logging(cfg_dict: dict | None):
|
|||||||
|
|
||||||
``cfg_dict`` is ``SyncConfig.logging`` (a dict or None). ``level`` is a
|
``cfg_dict`` is ``SyncConfig.logging`` (a dict or None). ``level`` is a
|
||||||
logging-level name string (default ``"INFO"``); ``path`` is the log file
|
logging-level name string (default ``"INFO"``); ``path`` is the log file
|
||||||
path (default ``"sync.log"``). The parent directory is created if missing.
|
path (default ``"sync.log"``). The parent directory is created if missing,
|
||||||
|
and any previously produced logs are archived before the new handler opens.
|
||||||
"""
|
"""
|
||||||
level = getattr(logging, (cfg_dict or {}).get("level", "INFO")) if cfg_dict else logging.INFO
|
level = getattr(logging, (cfg_dict or {}).get("level", "INFO")) if cfg_dict else logging.INFO
|
||||||
path = (cfg_dict or {}).get("path", "sync.log")
|
path = (cfg_dict or {}).get("path", "sync.log")
|
||||||
os.makedirs(os.path.dirname(path) or ".", exist_ok=True)
|
log_dir = os.path.dirname(path) or "."
|
||||||
|
os.makedirs(log_dir, exist_ok=True)
|
||||||
|
|
||||||
|
# Archive everything from prior runs so the root only shows today's log.
|
||||||
|
_archive_completed_logs(log_dir, os.path.basename(path))
|
||||||
|
|
||||||
|
# Idempotent: drop any handlers already attached to the root logger before
|
||||||
|
# re-adding. setup_logging can be called from more than one entry point
|
||||||
|
# (e.g. main.py and service.run), and the old code unconditionally
|
||||||
|
# addHandler'd each time, stacking duplicate handlers so every log line was
|
||||||
|
# written twice. Clearing first means repeated calls always yield exactly
|
||||||
|
# one file handler + one console handler, regardless of caller.
|
||||||
|
root = logging.getLogger()
|
||||||
|
for old in list(root.handlers):
|
||||||
|
root.removeHandler(old)
|
||||||
|
try:
|
||||||
|
old.close()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
root.setLevel(level)
|
||||||
|
|
||||||
|
cycle_filter = _CycleIdFilter()
|
||||||
|
|
||||||
h = logging.handlers.RotatingFileHandler(
|
h = logging.handlers.RotatingFileHandler(
|
||||||
path, maxBytes=10 * 1024 * 1024, backupCount=5, encoding="utf-8"
|
path, maxBytes=10 * 1024 * 1024, backupCount=5, encoding="utf-8"
|
||||||
)
|
)
|
||||||
h.setFormatter(logging.Formatter("%(asctime)s %(levelname)s [%(name)s] %(message)s"))
|
h.addFilter(cycle_filter)
|
||||||
root = logging.getLogger()
|
h.setFormatter(logging.Formatter(
|
||||||
root.setLevel(level)
|
"%(asctime)s %(levelname)s [%(name)s] [cyc:%(cycle)s] %(message)s"
|
||||||
|
))
|
||||||
root.addHandler(h)
|
root.addHandler(h)
|
||||||
|
|
||||||
sh = logging.StreamHandler()
|
sh = logging.StreamHandler()
|
||||||
|
sh.addFilter(cycle_filter)
|
||||||
|
# Console keeps the short format (fullsync/compare are interactive there).
|
||||||
sh.setFormatter(logging.Formatter("%(levelname)s %(message)s"))
|
sh.setFormatter(logging.Formatter("%(levelname)s %(message)s"))
|
||||||
root.addHandler(sh)
|
root.addHandler(sh)
|
||||||
|
|||||||
@@ -2,9 +2,24 @@
|
|||||||
|
|
||||||
``cycle(cfg)`` runs one full pass over every configured file:
|
``cycle(cfg)`` runs one full pass over every configured file:
|
||||||
1. capture — read each file's change log and stage rows into SyncQueue;
|
1. capture — read each file's change log and stage rows into SyncQueue;
|
||||||
2. apply — drain the queue via ``dbo.usp_SyncApply``;
|
2. apply — drain the queue via ``usp_SyncApply``;
|
||||||
3. cleanup — delete applied log rows from each file's Access log.
|
3. cleanup — delete applied log rows from each file's Access log.
|
||||||
|
|
||||||
|
Observability (rebuilt so data divergence is traceable from the log alone):
|
||||||
|
- every cycle gets a short correlation id; ``logging_setup`` stamps it on each
|
||||||
|
log line as ``[cyc:xxxxxxxx]`` and the same id is passed to ``usp_SyncApply``
|
||||||
|
so ``SyncApplyRunLog`` rows on SQL Server join back to the exact log lines of
|
||||||
|
the cycle that produced them;
|
||||||
|
- after apply, the per-table run-log rows (pending/merged/deleted/applied/
|
||||||
|
error/dead counts, duration, error message) are read back and logged --
|
||||||
|
the old single "apply done" line hid all of this;
|
||||||
|
- queue health is checked every cycle: ``error``/``dead`` rows, which
|
||||||
|
previously accumulated in complete silence, now emit WARNINGs with per-row
|
||||||
|
samples (table, record id, retry count, error message) -- these are exactly
|
||||||
|
the changes that exist in Access but never reached SQL Server;
|
||||||
|
- idle cycles log at DEBUG so INFO stays high-signal; ``run`` emits a periodic
|
||||||
|
idle heartbeat so a quiet log still proves the service is alive.
|
||||||
|
|
||||||
Each file's capture and cleanup is wrapped in its own try/except so one
|
Each file's capture and cleanup is wrapped in its own try/except so one
|
||||||
file's failure is logged and the cycle continues; the writer is always closed
|
file's failure is logged and the cycle continues; the writer is always closed
|
||||||
in a ``finally``. ``run(cfg)`` loops ``cycle`` with a sleep; ``main()``
|
in a ``finally``. ``run(cfg)`` loops ``cycle`` with a sleep; ``main()``
|
||||||
@@ -13,14 +28,15 @@ loads the config from ``argv[1]`` (default ``config.yaml``).
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
import sys
|
import sys
|
||||||
import time
|
import time
|
||||||
|
import uuid
|
||||||
import logging
|
import logging
|
||||||
|
|
||||||
from .config import load_config
|
from .config import load_config
|
||||||
from .access_reader import AccessReader
|
from .access_reader import AccessReader
|
||||||
from .sql_writer import SqlWriter
|
from .sql_writer import SqlWriter
|
||||||
from .capture import capture_file
|
from .capture import capture_file, CaptureStats
|
||||||
from .cleanup import cleanup_file
|
from .cleanup import cleanup_file
|
||||||
from .logging_setup import setup_logging
|
from .logging_setup import setup_logging, set_cycle_id
|
||||||
|
|
||||||
log = logging.getLogger("sync.service")
|
log = logging.getLogger("sync.service")
|
||||||
|
|
||||||
@@ -30,47 +46,156 @@ def run(cfg):
|
|||||||
|
|
||||||
Configures logging once on entry. Intended to be started by ``main()``
|
Configures logging once on entry. Intended to be started by ``main()``
|
||||||
under the service host (e.g. NSSM). Not unit-tested (infinite loop);
|
under the service host (e.g. NSSM). Not unit-tested (infinite loop);
|
||||||
``cycle()`` is the testable unit.
|
``cycle()`` is the testable unit. Emits an idle heartbeat every
|
||||||
|
``runtime.idle_heartbeat_seconds`` so a quiet log still proves liveness
|
||||||
|
now that idle cycles log at DEBUG.
|
||||||
"""
|
"""
|
||||||
setup_logging(cfg.logging)
|
setup_logging(cfg.logging)
|
||||||
|
log.info(
|
||||||
|
"service started: files=%d poll_interval=%ds idle_heartbeat=%ds",
|
||||||
|
len(cfg.files), cfg.runtime.poll_interval_seconds,
|
||||||
|
cfg.runtime.idle_heartbeat_seconds,
|
||||||
|
)
|
||||||
|
idle_since = None
|
||||||
|
idle_cycles = 0
|
||||||
while True:
|
while True:
|
||||||
cycle(cfg)
|
active = cycle(cfg)
|
||||||
|
now = time.monotonic()
|
||||||
|
if active:
|
||||||
|
idle_since, idle_cycles = None, 0
|
||||||
|
else:
|
||||||
|
idle_cycles += 1
|
||||||
|
if idle_since is None:
|
||||||
|
idle_since = now
|
||||||
|
elif now - idle_since >= cfg.runtime.idle_heartbeat_seconds:
|
||||||
|
log.info("idle heartbeat: %d cycles with no changes in the "
|
||||||
|
"last %ds", idle_cycles, int(now - idle_since))
|
||||||
|
idle_since, idle_cycles = now, 0
|
||||||
time.sleep(cfg.runtime.poll_interval_seconds)
|
time.sleep(cfg.runtime.poll_interval_seconds)
|
||||||
|
|
||||||
|
|
||||||
def cycle(cfg):
|
def cycle(cfg) -> bool:
|
||||||
"""One capture -> apply -> cleanup pass over all files.
|
"""One capture -> apply -> cleanup pass over all files.
|
||||||
|
|
||||||
Per-file capture/cleanup failures are logged and do not abort the cycle.
|
Returns True when the cycle did any work (captured / applied / cleaned /
|
||||||
Apply failure does not block cleanup. The writer is always closed in a
|
purged anything); ``run`` uses this for idle-heartbeat pacing. Per-file
|
||||||
|
capture/cleanup failures are logged and do not abort the cycle. Apply
|
||||||
|
failure does not block cleanup. The writer is always closed in a
|
||||||
``finally``. Safe to call directly from tests (does not sleep or loop).
|
``finally``. Safe to call directly from tests (does not sleep or loop).
|
||||||
"""
|
"""
|
||||||
writer = SqlWriter(cfg.sql_server.conn_str, cfg.sql_server.sync_queue_table)
|
cycle_id = uuid.uuid4().hex[:8]
|
||||||
|
set_cycle_id(cycle_id)
|
||||||
|
t0 = time.monotonic()
|
||||||
|
activity = False
|
||||||
|
writer = SqlWriter(
|
||||||
|
cfg.sql_server.conn_str,
|
||||||
|
cfg.sql_server.sync_queue_table,
|
||||||
|
cfg.sql_server.archive_table,
|
||||||
|
cfg.sql_server.apply_proc,
|
||||||
|
cfg.sql_server.apply_runlog_table,
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
total_captured = 0
|
# ---- capture ------------------------------------------------------
|
||||||
|
total = CaptureStats()
|
||||||
for fm in cfg.files:
|
for fm in cfg.files:
|
||||||
reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
|
reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
|
||||||
try:
|
try:
|
||||||
n = capture_file(fm, reader, writer, cfg)
|
total.merge(capture_file(fm, reader, writer, cfg))
|
||||||
total_captured += n
|
|
||||||
except Exception:
|
except Exception:
|
||||||
log.exception("capture failed for %s", fm.file)
|
log.exception("capture failed for %s", fm.file)
|
||||||
finally:
|
finally:
|
||||||
reader.close()
|
reader.close()
|
||||||
log.info("captured %d rows", total_captured)
|
if total.read:
|
||||||
|
activity = True
|
||||||
|
log.info(
|
||||||
|
"capture summary: read=%d enqueued=%d dedup_skipped=%d "
|
||||||
|
"downgraded=%d out_of_scope=%d unknown_op=%d ops={%s}",
|
||||||
|
total.read, total.enqueued, total.dedup_skipped,
|
||||||
|
total.downgraded, total.out_of_scope, total.unknown_op,
|
||||||
|
total.ops_str(),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
log.debug("capture summary: no new change-log rows in any file")
|
||||||
|
|
||||||
|
# ---- apply ----------------------------------------------------------
|
||||||
try:
|
try:
|
||||||
writer.call_apply(cfg.runtime.max_retries)
|
writer.call_apply(cfg.runtime.max_retries, cycle_id)
|
||||||
log.info("apply done")
|
runs = writer.apply_run_results(cycle_id)
|
||||||
|
if runs is None:
|
||||||
|
# Audit infra (SyncApplyRunLog / updated proc) not deployed
|
||||||
|
# yet: keep the legacy coarse line. sql_writer already warned
|
||||||
|
# once about what to deploy.
|
||||||
|
log.info("apply done")
|
||||||
|
elif not runs:
|
||||||
|
if total.enqueued:
|
||||||
|
# Rows were staged but the proc wrote no audit rows: the
|
||||||
|
# table exists but the deployed proc is probably the old
|
||||||
|
# version that does not write into it.
|
||||||
|
log.info("apply done (proc wrote no run-log rows -- "
|
||||||
|
"re-run the updated sql/02_sync_apply.sql?)")
|
||||||
|
else:
|
||||||
|
log.debug("apply done: queue was empty")
|
||||||
|
else:
|
||||||
|
activity = True
|
||||||
|
for r in runs:
|
||||||
|
if r["Outcome"] == "ok":
|
||||||
|
log.info(
|
||||||
|
"apply ok table=%s.%s pending=%s records=%s "
|
||||||
|
"merged=%s deleted=%s applied=%s dur_ms=%s",
|
||||||
|
r["TargetSchema"], r["TargetTable"],
|
||||||
|
r["PendingCount"], r["DistinctRecords"],
|
||||||
|
r["MergedCount"], r["DeletedCount"],
|
||||||
|
r["AppliedCount"], r["DurationMs"],
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
log.warning(
|
||||||
|
"apply FAILED table=%s.%s pending=%s records=%s "
|
||||||
|
"-> error=%s dead=%s dur_ms=%s msg=%s",
|
||||||
|
r["TargetSchema"], r["TargetTable"],
|
||||||
|
r["PendingCount"], r["DistinctRecords"],
|
||||||
|
r["ErrorCount"], r["DeadCount"],
|
||||||
|
r["DurationMs"], r["ErrorMsg"],
|
||||||
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
log.exception("apply failed")
|
log.exception("apply failed")
|
||||||
|
|
||||||
|
# ---- queue health ---------------------------------------------------
|
||||||
|
# error/dead rows are changes that exist in Access but never reached
|
||||||
|
# the SQL mirror. Before this check they accumulated with zero trace
|
||||||
|
# in this log -- the classic "data diverged, no idea why" scenario.
|
||||||
|
try:
|
||||||
|
status = writer.queue_status_summary()
|
||||||
|
stuck = status.get("error", 0) + status.get("dead", 0)
|
||||||
|
if stuck:
|
||||||
|
log.warning(
|
||||||
|
"queue health: %d stuck row(s) [%s] -- these changes are "
|
||||||
|
"NOT in SQL Server and will show up as data divergence",
|
||||||
|
stuck,
|
||||||
|
",".join(f"{k}={v}" for k, v in sorted(status.items())),
|
||||||
|
)
|
||||||
|
for s in writer.queue_error_samples(10):
|
||||||
|
log.warning(
|
||||||
|
" stuck row: table=%s.%s record_id=%s "
|
||||||
|
"source_log_id=%s op=%s status=%s retries=%s "
|
||||||
|
"captured_at=%s err=%s",
|
||||||
|
s["TargetSchema"], s["TargetTable"], s["RecordID"],
|
||||||
|
s["SourceLogID"], s["OperateType"], s["Status"],
|
||||||
|
s["RetryCount"], s["CapturedAt"], s["ErrorMsg"],
|
||||||
|
)
|
||||||
|
elif status.get("pending", 0):
|
||||||
|
log.warning(
|
||||||
|
"queue health: %d row(s) still pending after apply "
|
||||||
|
"(apply may have failed this cycle)", status["pending"],
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
log.exception("queue health check failed")
|
||||||
|
|
||||||
|
# ---- cleanup --------------------------------------------------------
|
||||||
for fm in cfg.files:
|
for fm in cfg.files:
|
||||||
reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
|
reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
|
||||||
try:
|
try:
|
||||||
c = cleanup_file(fm, reader, writer, cfg)
|
if cleanup_file(fm, reader, writer, cfg):
|
||||||
if c:
|
activity = True
|
||||||
log.info("cleaned %d log rows from %s", c, fm.file)
|
|
||||||
except Exception:
|
except Exception:
|
||||||
log.exception("cleanup failed for %s", fm.file)
|
log.exception("cleanup failed for %s", fm.file)
|
||||||
finally:
|
finally:
|
||||||
@@ -81,12 +206,18 @@ def cycle(cfg):
|
|||||||
try:
|
try:
|
||||||
purged = writer.purge_cleaned(cfg.runtime.cleaned_retention_hours)
|
purged = writer.purge_cleaned(cfg.runtime.cleaned_retention_hours)
|
||||||
if purged:
|
if purged:
|
||||||
log.info("purged %d cleaned queue rows", purged)
|
activity = True
|
||||||
|
log.info("purged %d cleaned queue rows (older than %dh)",
|
||||||
|
purged, cfg.runtime.cleaned_retention_hours)
|
||||||
except Exception:
|
except Exception:
|
||||||
log.exception("purge failed")
|
log.exception("purge failed")
|
||||||
|
|
||||||
|
(log.info if activity else log.debug)(
|
||||||
|
"cycle finished in %.2fs", time.monotonic() - t0)
|
||||||
|
return activity
|
||||||
finally:
|
finally:
|
||||||
writer.close()
|
writer.close()
|
||||||
|
set_cycle_id(None)
|
||||||
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
|
|||||||
@@ -1,26 +1,50 @@
|
|||||||
"""SQL-side writer for the Access -> SQL Server sync.
|
"""SQL-side writer for the Access -> SQL Server sync.
|
||||||
|
|
||||||
SqlWriter owns the pyodbc connection used to (a) dedup-insert rows into
|
SqlWriter owns the pyodbc connection used to (a) dedup-insert rows into the
|
||||||
``dbo.SyncQueue``, (b) invoke ``dbo.usp_SyncApply`` to drain the queue, and
|
staging queue (``ProductionDataBaseSync.SyncQueue`` by default), (b) invoke the
|
||||||
(c) report which SourceLogIDs have been applied.
|
apply proc (``ProductionDataBaseSync.usp_SyncApply``) to drain the queue,
|
||||||
|
(c) report which SourceLogIDs have been applied, and (d) append every consumed
|
||||||
|
Access change-log row to the permanent audit store
|
||||||
|
(``ProductionDataBaseSync.SyncLogArchive``). The queue/archive/proc names are
|
||||||
|
injected from config so the whole sync can live under a dedicated schema.
|
||||||
|
|
||||||
|
Observability additions:
|
||||||
|
- the dedup inserts now report whether a row was actually inserted (True) or
|
||||||
|
already present (False), so capture can log dedup hits -- the tell-tale of a
|
||||||
|
re-capture after a failed apply;
|
||||||
|
- ``call_apply`` passes the service's cycle correlation id to the proc, which
|
||||||
|
records one audit row per target table into
|
||||||
|
``ProductionDataBaseSync.SyncApplyRunLog`` (see sql/04_sync_apply_runlog.sql);
|
||||||
|
- ``apply_run_results`` reads those rows back so the service log shows
|
||||||
|
per-table merged/deleted/applied/error/dead counts instead of "apply done";
|
||||||
|
- ``queue_status_summary`` / ``queue_error_samples`` surface stuck
|
||||||
|
(``error``/``dead``) queue rows, which previously accumulated silently.
|
||||||
|
|
||||||
The connection is opened with ``autocommit=True`` on purpose. ``usp_SyncApply``
|
The connection is opened with ``autocommit=True`` on purpose. ``usp_SyncApply``
|
||||||
manages its own transaction internally (BEGIN TRAN ... ROLLBACK on error); if
|
manages its own transaction internally (BEGIN TRAN ... ROLLBACK on error); if
|
||||||
the caller held an outer implicit transaction, the proc's ROLLBACK would
|
the caller held an outer implicit transaction, the proc's ROLLBACK would
|
||||||
cascade and raise SQL error 266. The dedup ``IF NOT EXISTS ... INSERT`` is a
|
cascade and raise SQL error 266. The dedup ``INSERT .. SELECT .. WHERE NOT
|
||||||
single statement that is atomic under autocommit, so no explicit transaction is
|
EXISTS`` is a single statement that is atomic under autocommit, so no explicit
|
||||||
needed on the write path either.
|
transaction is needed on the write path either.
|
||||||
"""
|
"""
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
|
||||||
import pyodbc
|
import pyodbc
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# Process-wide deployment-state flags (SqlWriter instances are recreated every
|
||||||
|
# cycle, so per-instance flags would re-warn each cycle).
|
||||||
|
_proc_lacks_cycle_id = False # deployed usp_SyncApply predates @CycleID
|
||||||
|
_runlog_missing_noted = False # SyncApplyRunLog table not deployed yet
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class QueueRow:
|
class QueueRow:
|
||||||
"""A single staged change to enqueue into ``dbo.SyncQueue``."""
|
"""A single staged change to enqueue into the SyncQueue table."""
|
||||||
|
|
||||||
source_file: str
|
source_file: str
|
||||||
source_table: str
|
source_table: str
|
||||||
@@ -32,35 +56,67 @@ class QueueRow:
|
|||||||
row_data: str | None
|
row_data: str | None
|
||||||
|
|
||||||
|
|
||||||
class SqlWriter:
|
@dataclass
|
||||||
"""Writes to SyncQueue and drives the apply proc over a pyodbc connection."""
|
class ArchiveRow:
|
||||||
|
"""A single Access change-log row to append to ``SyncLogArchive``.
|
||||||
|
|
||||||
def __init__(self, conn_str: str, queue_table: str = "dbo.SyncQueue"):
|
Captures both the ORIGINAL operate type recorded by the Access data macro
|
||||||
|
and the PROCESSED operate type actually sent to the queue, so a downgrade
|
||||||
|
(e.g. Insert -> Delete when the row is momentarily unreadable) stays visible
|
||||||
|
forever. ``original_time`` preserves the Access log's own timestamp, which
|
||||||
|
the queue path discards.
|
||||||
|
"""
|
||||||
|
|
||||||
|
source_file: str
|
||||||
|
source_table: str
|
||||||
|
source_log_id: int
|
||||||
|
record_id: str
|
||||||
|
target_schema: str
|
||||||
|
target_table: str
|
||||||
|
original_operate_type: str
|
||||||
|
processed_operate_type: str
|
||||||
|
row_data: str | None
|
||||||
|
original_time: object
|
||||||
|
|
||||||
|
|
||||||
|
class SqlWriter:
|
||||||
|
"""Writes to SyncQueue/SyncLogArchive and drives the apply proc."""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
conn_str: str,
|
||||||
|
queue_table: str = "ProductionDataBaseSync.SyncQueue",
|
||||||
|
archive_table: str = "ProductionDataBaseSync.SyncLogArchive",
|
||||||
|
apply_proc: str = "ProductionDataBaseSync.usp_SyncApply",
|
||||||
|
runlog_table: str = "ProductionDataBaseSync.SyncApplyRunLog",
|
||||||
|
):
|
||||||
self.conn_str = conn_str
|
self.conn_str = conn_str
|
||||||
self.queue_table = queue_table
|
self.queue_table = queue_table
|
||||||
|
self.archive_table = archive_table
|
||||||
|
self.apply_proc = apply_proc
|
||||||
|
self.runlog_table = runlog_table
|
||||||
# autocommit=True: usp_SyncApply manages its own transaction internally.
|
# autocommit=True: usp_SyncApply manages its own transaction internally.
|
||||||
# An outer pyodbc transaction would conflict on ROLLBACK (SQL error 266).
|
# An outer pyodbc transaction would conflict on ROLLBACK (SQL error 266).
|
||||||
self._conn = pyodbc.connect(conn_str, autocommit=True)
|
self._conn = pyodbc.connect(conn_str, autocommit=True)
|
||||||
|
|
||||||
def insert_queue_row(self, row: QueueRow) -> None:
|
def insert_queue_row(self, row: QueueRow) -> bool:
|
||||||
"""Idempotently enqueue ``row`` (dedup on SourceFile/Table/LogID).
|
"""Idempotently enqueue ``row`` (dedup on SourceFile/Table/LogID).
|
||||||
|
|
||||||
``IF NOT EXISTS ... INSERT`` is a single statement, atomic under
|
Returns True when a new queue row was inserted, False on a dedup hit
|
||||||
autocommit. The unique index UX_SyncQueue_Dedup is the DB backstop.
|
(the row was already staged -- typically a re-capture after a failed
|
||||||
|
apply left the Access log row in place). ``INSERT .. SELECT .. WHERE
|
||||||
|
NOT EXISTS`` is a single statement, atomic under autocommit, and its
|
||||||
|
deterministic rowcount (0/1) is what makes the dedup outcome
|
||||||
|
observable for the capture audit log. The unique index
|
||||||
|
UX_SyncQueue_Dedup is the DB backstop.
|
||||||
"""
|
"""
|
||||||
# IF NOT EXISTS and the VALUES list each carry their own ? markers;
|
|
||||||
# pyodbc binds them positionally, so the 3 dedup keys are supplied
|
|
||||||
# twice (once for the EXISTS check, once for the INSERT).
|
|
||||||
cur = self._conn.cursor()
|
cur = self._conn.cursor()
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"IF NOT EXISTS (SELECT 1 FROM dbo.SyncQueue "
|
f"INSERT INTO {self.queue_table}(SourceFile,SourceTable,SourceLogID,"
|
||||||
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?) "
|
|
||||||
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,"
|
|
||||||
"TargetSchema,TargetTable,RecordID,OperateType,RowData,Status) "
|
"TargetSchema,TargetTable,RecordID,OperateType,RowData,Status) "
|
||||||
"VALUES (?,?,?,?,?,?,?,?, 'pending')",
|
"SELECT ?,?,?,?,?,?,?,?,'pending' "
|
||||||
row.source_file,
|
f"WHERE NOT EXISTS (SELECT 1 FROM {self.queue_table} "
|
||||||
row.source_table,
|
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?)",
|
||||||
row.source_log_id,
|
|
||||||
row.source_file,
|
row.source_file,
|
||||||
row.source_table,
|
row.source_table,
|
||||||
row.source_log_id,
|
row.source_log_id,
|
||||||
@@ -69,22 +125,147 @@ class SqlWriter:
|
|||||||
row.record_id,
|
row.record_id,
|
||||||
row.operate_type,
|
row.operate_type,
|
||||||
row.row_data,
|
row.row_data,
|
||||||
|
row.source_file,
|
||||||
|
row.source_table,
|
||||||
|
row.source_log_id,
|
||||||
)
|
)
|
||||||
# autocommit: statement already committed.
|
# autocommit: statement already committed.
|
||||||
|
return cur.rowcount > 0
|
||||||
|
|
||||||
def call_apply(self, max_retries: int) -> None:
|
def insert_archive_row(self, row: ArchiveRow) -> bool:
|
||||||
"""Drain the pending queue via the stored procedure.
|
"""Append ``row`` to the permanent audit store (dedup on source keys).
|
||||||
|
|
||||||
``usp_SyncApply`` flips rows to ``applied`` (or ``error`` after retries).
|
Written during capture, BEFORE cleanup deletes the Access log, so the
|
||||||
|
original evidence survives even after the queue row is purged. Stores
|
||||||
|
both the original and processed operate types plus the Access log time.
|
||||||
|
The ``WHERE NOT EXISTS`` guard keeps the first archive record if
|
||||||
|
capture re-runs the same log id (a prior cycle's apply failed and the
|
||||||
|
Access log persisted). Returns True when a new archive row was
|
||||||
|
written, False on a dedup hit.
|
||||||
"""
|
"""
|
||||||
cur = self._conn.cursor()
|
cur = self._conn.cursor()
|
||||||
cur.execute("EXEC dbo.usp_SyncApply ?", max_retries)
|
cur.execute(
|
||||||
|
f"INSERT INTO {self.archive_table}(SourceFile,SourceTable,SourceLogID,"
|
||||||
|
"RecordID,TargetSchema,TargetTable,OriginalOperateType,"
|
||||||
|
"ProcessedOperateType,RowData,OriginalTime) "
|
||||||
|
"SELECT ?,?,?,?,?,?,?,?,?,? "
|
||||||
|
f"WHERE NOT EXISTS (SELECT 1 FROM {self.archive_table} "
|
||||||
|
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?)",
|
||||||
|
row.source_file,
|
||||||
|
row.source_table,
|
||||||
|
row.source_log_id,
|
||||||
|
row.record_id,
|
||||||
|
row.target_schema,
|
||||||
|
row.target_table,
|
||||||
|
row.original_operate_type,
|
||||||
|
row.processed_operate_type,
|
||||||
|
row.row_data,
|
||||||
|
row.original_time,
|
||||||
|
row.source_file,
|
||||||
|
row.source_table,
|
||||||
|
row.source_log_id,
|
||||||
|
)
|
||||||
|
# autocommit: statement already committed.
|
||||||
|
return cur.rowcount > 0
|
||||||
|
|
||||||
|
def call_apply(self, max_retries: int, cycle_id: str | None = None) -> None:
|
||||||
|
"""Drain the pending queue via the stored procedure.
|
||||||
|
|
||||||
|
``usp_SyncApply`` flips rows to ``applied`` (or ``error``/``dead``
|
||||||
|
after retries) and -- once the updated proc plus the SyncApplyRunLog
|
||||||
|
table are deployed -- writes one audit row per target table tagged
|
||||||
|
with ``cycle_id``, so SQL-side apply stats join back to the service
|
||||||
|
log's ``[cyc:...]`` lines. Falls back to the legacy single-parameter
|
||||||
|
signature when the deployed proc predates ``@CycleID`` (SQL errors
|
||||||
|
8144/8145), so rolling out the Python side first keeps working.
|
||||||
|
"""
|
||||||
|
global _proc_lacks_cycle_id
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
if cycle_id is not None and not _proc_lacks_cycle_id:
|
||||||
|
try:
|
||||||
|
cur.execute(
|
||||||
|
f"EXEC {self.apply_proc} @MaxRetries=?, @CycleID=?",
|
||||||
|
max_retries, cycle_id,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
except pyodbc.Error as e:
|
||||||
|
msg = str(e)
|
||||||
|
if "8144" in msg or "8145" in msg:
|
||||||
|
_proc_lacks_cycle_id = True
|
||||||
|
log.warning(
|
||||||
|
"apply proc %s does not accept @CycleID yet -- run the "
|
||||||
|
"updated sql/02_sync_apply.sql to enable per-table "
|
||||||
|
"apply auditing; falling back to the legacy signature",
|
||||||
|
self.apply_proc,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
raise
|
||||||
|
cur.execute(f"EXEC {self.apply_proc} ?", max_retries)
|
||||||
|
|
||||||
|
def apply_run_results(self, cycle_id: str) -> list[dict] | None:
|
||||||
|
"""Per-table apply outcomes recorded by the proc for ``cycle_id``.
|
||||||
|
|
||||||
|
Reads ``SyncApplyRunLog`` (pending/distinct counts, MERGE/DELETE
|
||||||
|
rowcounts, applied/error/dead queue-row counts, duration and error
|
||||||
|
message per target table). Returns ``None`` when the audit
|
||||||
|
infrastructure is not deployed yet (legacy proc, or table missing) so
|
||||||
|
the caller can fall back to coarse logging; returns ``[]`` when the
|
||||||
|
queue was simply empty.
|
||||||
|
"""
|
||||||
|
global _runlog_missing_noted
|
||||||
|
if _proc_lacks_cycle_id:
|
||||||
|
return None # legacy proc never writes run-log rows
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
try:
|
||||||
|
cur.execute(
|
||||||
|
"SELECT TargetSchema,TargetTable,PendingCount,DistinctRecords,"
|
||||||
|
"MergedCount,DeletedCount,AppliedCount,ErrorCount,DeadCount,"
|
||||||
|
"Outcome,ErrorMsg,StartedAt,DurationMs "
|
||||||
|
f"FROM {self.runlog_table} WHERE CycleID=? ORDER BY RunLogID",
|
||||||
|
cycle_id,
|
||||||
|
)
|
||||||
|
except pyodbc.Error:
|
||||||
|
if not _runlog_missing_noted:
|
||||||
|
_runlog_missing_noted = True
|
||||||
|
log.warning(
|
||||||
|
"run-log table %s not available -- run "
|
||||||
|
"sql/04_sync_apply_runlog.sql to enable per-table apply "
|
||||||
|
"stats in this log", self.runlog_table,
|
||||||
|
)
|
||||||
|
return None
|
||||||
|
cols = [c[0] for c in cur.description]
|
||||||
|
return [dict(zip(cols, r)) for r in cur.fetchall()]
|
||||||
|
|
||||||
|
def queue_status_summary(self) -> dict[str, int]:
|
||||||
|
"""Row counts per Status (pending/applied/error/dead/cleaned)."""
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
cur.execute(
|
||||||
|
f"SELECT Status, COUNT(*) FROM {self.queue_table} GROUP BY Status"
|
||||||
|
)
|
||||||
|
return {r[0]: r[1] for r in cur.fetchall()}
|
||||||
|
|
||||||
|
def queue_error_samples(self, limit: int = 10) -> list[dict]:
|
||||||
|
"""Most recent ``error``/``dead`` queue rows, for WARNING-level triage.
|
||||||
|
|
||||||
|
These are exactly the changes that exist in Access but never reached
|
||||||
|
the SQL mirror -- the prime suspects for any data divergence.
|
||||||
|
"""
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
cur.execute(
|
||||||
|
"SELECT TOP (?) TargetSchema,TargetTable,RecordID,SourceLogID,"
|
||||||
|
"OperateType,Status,RetryCount,ErrorMsg,CapturedAt "
|
||||||
|
f"FROM {self.queue_table} WHERE Status IN ('error','dead') "
|
||||||
|
"ORDER BY QueueID DESC",
|
||||||
|
limit,
|
||||||
|
)
|
||||||
|
cols = [c[0] for c in cur.description]
|
||||||
|
return [dict(zip(cols, r)) for r in cur.fetchall()]
|
||||||
|
|
||||||
def applied_log_ids(self, source_file: str) -> list[int]:
|
def applied_log_ids(self, source_file: str) -> list[int]:
|
||||||
"""Return applied SourceLogIDs for ``source_file`` in ascending order."""
|
"""Return applied SourceLogIDs for ``source_file`` in ascending order."""
|
||||||
cur = self._conn.cursor()
|
cur = self._conn.cursor()
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"SELECT SourceLogID FROM dbo.SyncQueue "
|
f"SELECT SourceLogID FROM {self.queue_table} "
|
||||||
"WHERE SourceFile=? AND Status='applied' ORDER BY SourceLogID",
|
"WHERE SourceFile=? AND Status='applied' ORDER BY SourceLogID",
|
||||||
source_file,
|
source_file,
|
||||||
)
|
)
|
||||||
@@ -107,7 +288,7 @@ class SqlWriter:
|
|||||||
chunk = source_log_ids[i:i + 1000]
|
chunk = source_log_ids[i:i + 1000]
|
||||||
placeholders = ",".join("?" * len(chunk))
|
placeholders = ",".join("?" * len(chunk))
|
||||||
cur.execute(
|
cur.execute(
|
||||||
f"UPDATE dbo.SyncQueue SET Status='cleaned', "
|
f"UPDATE {self.queue_table} SET Status='cleaned', "
|
||||||
f"CleanedAt=sysdatetime() "
|
f"CleanedAt=sysdatetime() "
|
||||||
f"WHERE SourceFile=? AND Status='applied' "
|
f"WHERE SourceFile=? AND Status='applied' "
|
||||||
f"AND SourceLogID IN ({placeholders})",
|
f"AND SourceLogID IN ({placeholders})",
|
||||||
@@ -120,12 +301,12 @@ class SqlWriter:
|
|||||||
|
|
||||||
Keeps the table bounded: cleanup marks rows ``cleaned`` every cycle,
|
Keeps the table bounded: cleanup marks rows ``cleaned`` every cycle,
|
||||||
and this removes the old ones after a short audit/debug window so
|
and this removes the old ones after a short audit/debug window so
|
||||||
``dbo.SyncQueue`` stops growing without bound. Returns the number of
|
the queue table stops growing without bound. Returns the number of
|
||||||
rows removed.
|
rows removed.
|
||||||
"""
|
"""
|
||||||
cur = self._conn.cursor()
|
cur = self._conn.cursor()
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"DELETE FROM dbo.SyncQueue "
|
f"DELETE FROM {self.queue_table} "
|
||||||
"WHERE Status='cleaned' "
|
"WHERE Status='cleaned' "
|
||||||
"AND CleanedAt < DATEADD(hour, -?, GETDATE())",
|
"AND CleanedAt < DATEADD(hour, -?, GETDATE())",
|
||||||
retention_hours,
|
retention_hours,
|
||||||
|
|||||||
@@ -1,15 +1,23 @@
|
|||||||
"""Integration test for dbo.usp_SyncApply.
|
"""Integration test for usp_SyncApply.
|
||||||
|
|
||||||
Validates: IDENTITY-preserving INSERT, last-write-wins UPDATE, BIT conversion
|
Validates: IDENTITY-preserving INSERT, last-write-wins UPDATE, BIT conversion
|
||||||
from JSON ``true``/``false``, and DELETE of the last op. Creates a throwaway
|
from JSON ``true``/``false``, and DELETE of the last op. Creates a throwaway
|
||||||
schema ``sync_test`` and table ``ApplyDemo_YEAR2026`` and cleans them up at the
|
schema ``sync_test`` and table ``ApplyDemo_YEAR2026`` and cleans them up at the
|
||||||
end so no residue is left on CompanyDB.
|
end so no residue is left on CompanyDB.
|
||||||
|
|
||||||
|
Queue table and apply proc names are read from ``config.yaml`` so the test
|
||||||
|
follows whichever schema the deployment targets (currently ProductionDataBaseSync).
|
||||||
"""
|
"""
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
|
from sync.config import load_config
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.integration
|
@pytest.mark.integration
|
||||||
def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
||||||
|
cfg = load_config("config.yaml")
|
||||||
|
qt = cfg.sql_server.sync_queue_table
|
||||||
|
proc = cfg.sql_server.apply_proc
|
||||||
cur = sql_conn.cursor()
|
cur = sql_conn.cursor()
|
||||||
sch, tbl = "sync_test", "ApplyDemo_YEAR2026"
|
sch, tbl = "sync_test", "ApplyDemo_YEAR2026"
|
||||||
|
|
||||||
@@ -27,17 +35,17 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
|||||||
"时间 DATETIME2 NULL, "
|
"时间 DATETIME2 NULL, "
|
||||||
"标记 BIT NULL)"
|
"标记 BIT NULL)"
|
||||||
)
|
)
|
||||||
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'")
|
cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
|
||||||
|
|
||||||
# Insert then a later Update for the same RecordID=1 -> last write wins.
|
# Insert then a later Update for the same RecordID=1 -> last write wins.
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
"INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
||||||
"TargetTable,RecordID,OperateType,RowData,Status) "
|
"TargetTable,RecordID,OperateType,RowData,Status) "
|
||||||
"VALUES('t.accdb','ApplyDemo',1,'sync_test','ApplyDemo_YEAR2026','1',"
|
"VALUES('t.accdb','ApplyDemo',1,'sync_test','ApplyDemo_YEAR2026','1',"
|
||||||
"'Insert','{\"ID\":1,\"名字\":\"A\",\"数量\":3,\"时间\":\"2026-01-01T00:00:00\",\"标记\":true}','pending')"
|
"'Insert','{\"ID\":1,\"名字\":\"A\",\"数量\":3,\"时间\":\"2026-01-01T00:00:00\",\"标记\":true}','pending')"
|
||||||
)
|
)
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
"INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
||||||
"TargetTable,RecordID,OperateType,RowData,Status) "
|
"TargetTable,RecordID,OperateType,RowData,Status) "
|
||||||
"VALUES('t.accdb','ApplyDemo',2,'sync_test','ApplyDemo_YEAR2026','1',"
|
"VALUES('t.accdb','ApplyDemo',2,'sync_test','ApplyDemo_YEAR2026','1',"
|
||||||
"'Update','{\"ID\":1,\"名字\":\"A2\",\"数量\":5,\"时间\":\"2026-01-02T00:00:00\",\"标记\":false}','pending')"
|
"'Update','{\"ID\":1,\"名字\":\"A2\",\"数量\":5,\"时间\":\"2026-01-02T00:00:00\",\"标记\":false}','pending')"
|
||||||
@@ -45,7 +53,7 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
|||||||
sql_conn.commit()
|
sql_conn.commit()
|
||||||
|
|
||||||
# --- Act 1: apply upsert ----------------------------------------------
|
# --- Act 1: apply upsert ----------------------------------------------
|
||||||
cur.execute("EXEC dbo.usp_SyncApply @MaxRetries=5")
|
cur.execute("EXEC " + proc + " @MaxRetries=5")
|
||||||
sql_conn.commit()
|
sql_conn.commit()
|
||||||
|
|
||||||
# --- Assert 1: the later Update wins; BIT false -> 0 ------------------
|
# --- Assert 1: the later Update wins; BIT false -> 0 ------------------
|
||||||
@@ -57,15 +65,15 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
|||||||
assert row.标记 == 0 # BIT false
|
assert row.标记 == 0 # BIT false
|
||||||
|
|
||||||
# --- Act 2: a later Delete wins ---------------------------------------
|
# --- Act 2: a later Delete wins ---------------------------------------
|
||||||
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'")
|
cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
"INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
|
||||||
"TargetTable,RecordID,OperateType,RowData,Status) "
|
"TargetTable,RecordID,OperateType,RowData,Status) "
|
||||||
"VALUES('t.accdb','ApplyDemo',3,'sync_test','ApplyDemo_YEAR2026','1',"
|
"VALUES('t.accdb','ApplyDemo',3,'sync_test','ApplyDemo_YEAR2026','1',"
|
||||||
"'Delete',NULL,'pending')"
|
"'Delete',NULL,'pending')"
|
||||||
)
|
)
|
||||||
sql_conn.commit()
|
sql_conn.commit()
|
||||||
cur.execute("EXEC dbo.usp_SyncApply @MaxRetries=5")
|
cur.execute("EXEC " + proc + " @MaxRetries=5")
|
||||||
sql_conn.commit()
|
sql_conn.commit()
|
||||||
|
|
||||||
# --- Assert 2: row removed --------------------------------------------
|
# --- Assert 2: row removed --------------------------------------------
|
||||||
@@ -77,5 +85,5 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
|
|||||||
"IF OBJECT_ID('sync_test.ApplyDemo_YEAR2026') IS NOT NULL "
|
"IF OBJECT_ID('sync_test.ApplyDemo_YEAR2026') IS NOT NULL "
|
||||||
"DROP TABLE sync_test.ApplyDemo_YEAR2026"
|
"DROP TABLE sync_test.ApplyDemo_YEAR2026"
|
||||||
)
|
)
|
||||||
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'")
|
cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
|
||||||
sql_conn.commit()
|
sql_conn.commit()
|
||||||
|
|||||||
@@ -21,10 +21,16 @@ def test_insert_dedup_and_applied_ids():
|
|||||||
if not os.environ.get("RUN_INTEGRATION"):
|
if not os.environ.get("RUN_INTEGRATION"):
|
||||||
pytest.skip("integration")
|
pytest.skip("integration")
|
||||||
cfg = load_config("config.yaml")
|
cfg = load_config("config.yaml")
|
||||||
w = SqlWriter(cfg.sql_server.conn_str, "dbo.SyncQueue")
|
qt = cfg.sql_server.sync_queue_table
|
||||||
|
w = SqlWriter(
|
||||||
|
cfg.sql_server.conn_str,
|
||||||
|
qt,
|
||||||
|
cfg.sql_server.archive_table,
|
||||||
|
cfg.sql_server.apply_proc,
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
cur = w._conn.cursor()
|
cur = w._conn.cursor()
|
||||||
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'")
|
cur.execute(f"DELETE {qt} WHERE SourceFile='sqlw_test.accdb'")
|
||||||
qr = QueueRow(
|
qr = QueueRow(
|
||||||
source_file="sqlw_test.accdb",
|
source_file="sqlw_test.accdb",
|
||||||
source_table="T",
|
source_table="T",
|
||||||
@@ -38,7 +44,7 @@ def test_insert_dedup_and_applied_ids():
|
|||||||
w.insert_queue_row(qr)
|
w.insert_queue_row(qr)
|
||||||
w.insert_queue_row(qr) # duplicate must be deduped (ignored)
|
w.insert_queue_row(qr) # duplicate must be deduped (ignored)
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"SELECT COUNT(*) FROM dbo.SyncQueue "
|
f"SELECT COUNT(*) FROM {qt} "
|
||||||
"WHERE SourceFile='sqlw_test.accdb' AND SourceLogID=100"
|
"WHERE SourceFile='sqlw_test.accdb' AND SourceLogID=100"
|
||||||
)
|
)
|
||||||
assert cur.fetchone()[0] == 1
|
assert cur.fetchone()[0] == 1
|
||||||
@@ -47,12 +53,12 @@ def test_insert_dedup_and_applied_ids():
|
|||||||
w.call_apply(max_retries=5)
|
w.call_apply(max_retries=5)
|
||||||
|
|
||||||
cur.execute(
|
cur.execute(
|
||||||
"UPDATE dbo.SyncQueue SET Status='applied' "
|
f"UPDATE {qt} SET Status='applied' "
|
||||||
"WHERE SourceFile='sqlw_test.accdb'"
|
"WHERE SourceFile='sqlw_test.accdb'"
|
||||||
)
|
)
|
||||||
assert w.applied_log_ids("sqlw_test.accdb") == [100]
|
assert w.applied_log_ids("sqlw_test.accdb") == [100]
|
||||||
finally:
|
finally:
|
||||||
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'")
|
cur.execute(f"DELETE {qt} WHERE SourceFile='sqlw_test.accdb'")
|
||||||
w.close()
|
w.close()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user