Merge feat/auto-create-table-and-verify-sync into main

Merge auto table creation with per-PK sync verification, combining
both features: target tables are auto-created when missing, and
every sync batch is verified by primary key before marking synced.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Misaka_Company
2026-06-12 10:52:54 +08:00
5 changed files with 329 additions and 111 deletions

View File

@@ -45,7 +45,7 @@ EXCEL_SYNC_UPTIME_KUMA_CONFIG = {
# ================= 运行参数 =================
# 合并原 config.py 和 update_config.py 的配置
POLL_INTERVAL = 5 # 轮询间隔(秒)
POLL_INTERVAL = 30 # 轮询间隔(秒)
EXCEL_SYNC_INTERVAL = 600 # Excel 同步周期(秒),默认 10 分钟
BATCH_SIZE = 10000 # 批量处理大小
CACHE_DIR = os.path.join(os.getcwd(), "temp") # Excel 缓存目录

View File

@@ -26,6 +26,16 @@ SYNC_MAPPING = {
"target_schema": "productionContractData",
"target_table": "25年温度计合同数据",
"pk_col": "ID"
},
"25年变送器合同数据": {
"target_schema": "productionContractData",
"target_table": "25年变送器合同数据",
"pk_col": "ID"
},
"26年变送器合同数据": {
"target_schema": "productionContractData",
"target_table": "26年变送器合同数据",
"pk_col": "ID"
}
},
r"\\192.168.110.114\生产进度表\2026年数据\成品入库.accdb": {
@@ -301,6 +311,13 @@ SYNC_MAPPING = {
"target_table": "执行卡下发记录_YEAR2026",
"pk_col": "ID"
}
},
r"\\192.168.110.114\生产进度表\2026年数据\弯管车间.accdb": {
"执行卡下发记录": {
"target_schema": "tubeBending",
"target_table": "烘洗_YEAR2026",
"pk_col": "ID"
}
}
}

View File

@@ -1,7 +1,82 @@
# db_utils.py
import pyodbc
import datetime
import decimal
from config import SQL_SERVER_CONN, ACCESS_DRIVER
# ================= 表存在性检查 & 自动建表 =================
def table_exists(cursor, schema, table):
"""检查 SQL Server 表是否存在"""
query = """
SELECT COUNT(*)
FROM sys.tables t
JOIN sys.schemas s ON t.schema_id = s.schema_id
WHERE s.name = ? AND t.name = ?
"""
cursor.execute(query, (schema, table))
return cursor.fetchone()[0] > 0
def ensure_schema(cursor, schema):
"""确保 schema 存在,不存在则创建"""
cursor.execute("SELECT SCHEMA_ID(?)", (schema,))
if cursor.fetchone()[0] is None:
cursor.execute(f"CREATE SCHEMA [{schema}]")
def _access_col_to_sql(col_info, pk_col):
"""将 Access 列描述 (cursor.description 元组) 转为 SQL Server 列定义
Access ODBC 驱动的 type_code 是 Python 类型对象:
int → INT
str → NVARCHAR(size) (size > 4000 时为 Memo 字段 → NVARCHAR(MAX))
datetime.datetime → DATETIME
float → FLOAT
decimal.Decimal → DECIMAL(p,s)
bool → BIT
"""
col_name, type_code, _, size, precision, scale, nullable = col_info
is_pk = (col_name == pk_col)
if type_code is int:
sql_type = "INT"
if is_pk:
sql_type += " IDENTITY(1,1) PRIMARY KEY"
elif type_code is float:
sql_type = "FLOAT"
elif type_code is bool:
sql_type = "BIT"
elif type_code is datetime.datetime:
sql_type = "DATETIME"
elif type_code is decimal.Decimal:
sql_type = f"DECIMAL({precision or 18}, {scale or 0})"
elif type_code is str:
# size > 4000 → Access Memo 字段,用 NVARCHAR(MAX)
if not size or size <= 0 or size > 4000:
sql_type = "NVARCHAR(MAX)"
else:
sql_type = f"NVARCHAR({size})"
if is_pk:
sql_type += " PRIMARY KEY"
else:
sql_type = "NVARCHAR(255)"
if is_pk:
sql_type += " PRIMARY KEY"
return f"[{col_name}] {sql_type}"
def create_table_from_access(sql_cursor, target_schema, target_table,
acc_description, pk_col):
"""根据 Access cursor.description 在 SQL Server 自动建表"""
col_defs = [_access_col_to_sql(col, pk_col) for col in acc_description]
full_name = fmt_table(target_schema, target_table)
col_str = ",\n ".join(col_defs)
create_sql = f"CREATE TABLE {full_name} (\n {col_str}\n)"
sql_cursor.execute(create_sql)
def get_sql_conn():
"""获取 SQL Server 连接"""
# 显式添加 TrustServerCertificate=yes 以兼容 ODBC Driver 18+

View File

@@ -140,7 +140,16 @@ def run_full_sync():
log_info( f" 检测到 {len(columns)} 个列")
insert_sql = db_utils.generate_insert_sql(target_schema, target_table, columns)
# 3. TRUNCATE 目标表
# 3. 检查目标表是否存在,不存在则自动创建
if not db_utils.table_exists(sql_cursor, target_schema, target_table):
db_utils.ensure_schema(sql_cursor, target_schema)
db_utils.create_table_from_access(
sql_cursor, target_schema, target_table,
acc_cursor.description, target_config['pk_col']
)
sql_conn.commit() # DDL 必须单独提交,否则 pyodbc 参数绑定无法获取正确的列元数据
log_info( f" 目标表不存在,已自动创建")
else:
sql_cursor.execute(f"TRUNCATE TABLE {full_target_name}")
log_info( f" 已清空目标表")

View File

@@ -15,6 +15,18 @@ LoggerManager("run_incremental_sync", log_prefix="incremental")
uptime_monitor = UptimeKumaMonitor(UPTIME_KUMA_CONFIG)
uptime_monitor.set_logger(log_warning)
# 提交后是否再做一次独立复核 (额外保险, 略增开销; 若不需要可设为 False)
ENABLE_POST_COMMIT_VERIFY = True
# IN 子句单批最大参数数 (SQL Server 上限约 2100, 留足余量)
IN_CLAUSE_BATCH = 900
class SyncVerificationError(Exception):
"""数据落库校验未通过时抛出, 触发事务回滚并保留 Synced=0 以便下一轮重试"""
pass
# ================= 辅助函数 =================
def has_identity_column(sql_cursor, schema, table):
@@ -29,12 +41,90 @@ def has_identity_column(sql_cursor, schema, table):
sql_cursor.execute(query, (schema, table))
return sql_cursor.fetchone()[0] > 0
def _chunked(values, size=IN_CLAUSE_BATCH):
"""将集合/列表分批, 避免 IN 子句参数数量超过 SQL Server 上限"""
items = list(values)
for i in range(0, len(items), size):
yield items[i:i + size]
def _norm_key(v):
"""主键归一化, 用于跨来源(日志侧 / Access 侧 / SQL 侧)的集合比较, 规避类型差异"""
return str(v).strip() if v is not None else None
def fetch_existing_pks(sql_cursor, full_name, pk_col, pk_values):
"""
拿一批主键去目标表查询, 返回其中【实际存在】的主键集合 (分批查询)。
这是逐主键校验的基础: 用 "查到了哪些" 来精确判断存在 / 不存在,
而不是用数量相减 (数量法会被 "一边少插一边漏删" 互相抵消而误判)。
"""
found = set()
for chunk in _chunked(pk_values):
if not chunk:
continue
placeholders = ','.join(['?'] * len(chunk))
sql = f"SELECT [{pk_col}] FROM {full_name} WHERE [{pk_col}] IN ({placeholders})"
sql_cursor.execute(sql, list(chunk))
for row in sql_cursor.fetchall():
found.add(row[0])
return found
def verify_sync_result(sql_cursor, full_name, pk_col, record_ids, inserted_pk_set):
"""
逐主键校验本批同步结果 (不依赖数量比较):
- 删除校验: 源端已读不到的记录 (record_ids 中不在 inserted_pk_set 的部分),
删除后必须【全部不存在】于目标表; 只要还查得到任意一条即判失败。
- 插入校验: 从 Access 实际读到并重新插入的主键 (inserted_pk_set),
必须【全部存在】于目标表; 只要有一条查不到即判失败。
校验不通过抛出 SyncVerificationError, 并给出具体出问题的主键样例。
返回值: 本批 "随源删除" 的记录条数 (供日志展示)。
"""
inserted_keys = {_norm_key(pk) for pk in inserted_pk_set}
# 1) 删除校验: 应删除的记录, 删除后必须不在目标表
to_delete = [rid for rid in record_ids if _norm_key(rid) not in inserted_keys]
if to_delete:
leftover = fetch_existing_pks(sql_cursor, full_name, pk_col, to_delete)
if leftover:
sample = list(leftover)[:5]
raise SyncVerificationError(
f"删除校验失败: {len(leftover)} 条记录删除后仍存在于目标表, 例如 {sample}")
# 2) 插入校验: 应插入的记录, 插入后必须存在于目标表
if inserted_pk_set:
present_keys = {_norm_key(p) for p in fetch_existing_pks(sql_cursor, full_name, pk_col, inserted_pk_set)}
missing = [pk for pk in inserted_pk_set if _norm_key(pk) not in present_keys]
if missing:
raise SyncVerificationError(
f"插入校验失败: {len(missing)} 条记录插入后在目标表查不到, 例如 {missing[:5]}")
return len(to_delete)
def find_pk_index(acc_cols, pk_col):
"""在 Access 结果列中定位主键列下标 (大小写不敏感兜底)"""
if pk_col in acc_cols:
return acc_cols.index(pk_col)
lower_map = {c.lower(): i for i, c in enumerate(acc_cols)}
return lower_map.get(pk_col.lower())
# ================= 主逻辑 =================
def process_sync_task():
"""
逻辑重构
逻辑重构:
遍历 SYNC_MAPPING 中的每一个文件 -> 去日志表查询该文件下特定表的未同步记录。
校验机制 (逐主键确认):
删旧插新后, 在提交前对本批数据做删除校验 + 插入校验 ——
应删除的主键确认已不在目标表、应插入的主键确认已在目标表, 二者都通过才标记
Synced=1 并提交; 否则回滚整批, 保持 Synced=0, 下一轮自动重试。
"""
sql_conn = db_utils.get_sql_conn()
sql_cursor = sql_conn.cursor()
@@ -51,25 +141,21 @@ def process_sync_task():
for clean_path, tables_map in SYNC_MAPPING.items():
# 1. 准备查询条件
# 获取该文件下所有需要同步的表名列表
target_tables = list(tables_map.keys())
if not target_tables:
continue
# 构造 TableAddress 的精确匹配条件
# VBA 逻辑:网络路径带 ";DATABASE=", 本地路径带 "LOCAL:"
# 我们直接构造这两个字符串,让 SQL Server 做精确匹配,效率极高
# 注意:将路径转为 Windows 标准反斜杠
win_path = os.path.normpath(clean_path)
addr_candidates = [
f";DATABASE={win_path}", # 情况1
f"LOCAL={win_path}", # 情况2 (注意 VBA 代码里可能是 LOCAL: 或 LOCAL=,请核对)
f"LOCAL={win_path}", # 情况2
f"LOCAL:{win_path}", # 情况3
win_path # 情况4 (兼容没有前缀的情况)
]
# 2. 构造动态 SQL 查询
# WHERE Synced=0 AND Address IN (...) AND TableName IN (...)
placeholders_addr = ','.join(['?'] * len(addr_candidates))
placeholders_tbl = ','.join(['?'] * len(target_tables))
@@ -99,7 +185,7 @@ def process_sync_task():
log_file(f"{os.path.basename(clean_path)} 发现 {len(logs)} 条待同步变更")
# 3. 本地分组 (按表名)
# 结构: table_tasks[TableName] = { ids: {}, log_ids: [] }
# 结构: table_tasks[TableName] = { record_ids: set(), log_ids: [] }
table_tasks = {}
for row in logs:
log_id, acc_table, record_id = row
@@ -127,6 +213,7 @@ def process_sync_task():
file_error_count = 0
file_total_records = 0
try:
for acc_table, data in table_tasks.items():
record_ids = list(data['record_ids'])
log_ids = data['log_ids']
@@ -138,8 +225,6 @@ def process_sync_task():
pk_col = target_conf['pk_col']
target_full_name = db_utils.fmt_table(target_schema, target_table)
log_processing(f"正在同步表 [{acc_table}] → [{target_table}] ({len(record_ids)} 条记录)")
# 检查目标表是否存在,不存在则自动创建
if not db_utils.table_exists(sql_cursor, target_schema, target_table):
db_utils.ensure_schema(sql_cursor, target_schema)
@@ -150,46 +235,54 @@ def process_sync_task():
)
log_info(f"目标表 [{target_table}] 不存在,已自动创建")
# 事务状态追踪
log_processing(f"正在同步表 [{acc_table}] → [{target_table}] ({len(record_ids)} 条记录)")
identity_enabled = False
transaction_begun = False
try:
# --- 开始事务 ---
sql_cursor.execute("BEGIN TRANSACTION")
transaction_begun = True
# --- A. Access 查新数据 ---
# --- A. 从 Access 读取最新数据 ---
ids_placeholders = ','.join(['?'] * len(record_ids))
acc_sql = f"SELECT * FROM [{acc_table}] WHERE [{pk_col}] IN ({ids_placeholders})"
acc_cursor.execute(acc_sql, record_ids)
new_rows = acc_cursor.fetchall()
acc_cols = [col[0] for col in acc_cursor.description]
# --- B. SQL Server 删旧插新 ---
# 1. 删除旧数据
# 定位主键列, 取出 Access 实际读到的主键集合 (校验基准)
pk_index = find_pk_index(acc_cols, pk_col)
if new_rows and pk_index is None:
raise SyncVerificationError(
f"在 Access 表 [{acc_table}] 中找不到主键列 [{pk_col}]")
inserted_pk_set = (set(row[pk_index] for row in new_rows)
if pk_index is not None else set())
# --- B. SQL Server 删除旧记录 ---
del_sql = f"DELETE FROM {target_full_name} WHERE [{pk_col}] IN ({ids_placeholders})"
sql_cursor.execute(del_sql, record_ids)
# 2. 插入新数据
# --- C. 插入新记录 ---
if new_rows:
insert_sql = db_utils.generate_insert_sql(target_schema, target_table, acc_cols)
# 检查并启用 IDENTITY_INSERT
has_identity = has_identity_column(sql_cursor, target_schema, target_table)
if has_identity:
if has_identity_column(sql_cursor, target_schema, target_table):
try:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} ON")
identity_enabled = True
except Exception as id_err:
log_error(f"无法启用 IDENTITY_INSERT: {id_err}")
raise
sql_cursor.executemany(insert_sql, new_rows)
# 关闭 IDENTITY_INSERT
# 立即关闭 IDENTITY_INSERT (会话级设置, 不随事务回滚)
if identity_enabled:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
identity_enabled = False
# 3. 标记日志 Synced = 1
# --- D. 提交前校验 (逐主键确认: 应删的已不在 / 应插的已在) ---
removed_count = verify_sync_result(
sql_cursor, target_full_name, pk_col, record_ids, inserted_pk_set)
# --- E. 校验通过 -> 标记日志 Synced = 1 ---
log_placeholders = ','.join(['?'] * len(log_ids))
update_log_sql = f"""
UPDATE {log_full_name}
@@ -198,30 +291,51 @@ def process_sync_task():
"""
sql_cursor.execute(update_log_sql, log_ids)
# 提交事务
# --- F. 提交事务 (删/插/标记 原子生效) ---
sql_conn.commit()
transaction_begun = False
log_info(f"表 [{target_table}] 同步完成: {len(record_ids)} 条记录")
# --- G. 提交后独立复核 (可选, 防 "提交成功但未持久化" 等极端情况) ---
if ENABLE_POST_COMMIT_VERIFY:
try:
verify_sync_result(
sql_cursor, target_full_name, pk_col, record_ids, inserted_pk_set)
except SyncVerificationError as post_err:
log_critical(
f"严重: 表 [{target_table}] 提交后复核失败! {post_err}; "
f"撤销同步标记以便重试")
revert_sql = f"""
UPDATE {log_full_name}
SET {cols['col_synced']} = 0
WHERE {cols['col_log_id']} IN ({log_placeholders})
"""
sql_cursor.execute(revert_sql, log_ids)
sql_conn.commit()
file_error_count += 1
continue
# 同步成功
msg = f"表 [{target_table}] 同步并校验通过: {len(inserted_pk_set)} 条入库"
if removed_count > 0:
msg += f", {removed_count} 条随源删除"
log_success(msg)
file_success_count += 1
file_total_records += len(record_ids)
except Exception as tbl_err:
# 回滚事务
if transaction_begun:
sql_conn.rollback()
transaction_begun = False
# 清理 IDENTITY_INSERT 状态
# 清理可能残留的 IDENTITY_INSERT 会话状态
if identity_enabled:
try:
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
except:
pass
log_error(f"表 [{acc_table}] 同步失败: {tbl_err}")
# 回滚本表的 删/插/标记, 保持 Synced=0, 下一轮自动重试
try:
sql_conn.rollback()
except:
pass
log_error(f"表 [{acc_table}] 同步失败 (已回滚, 将重试): {tbl_err}")
file_error_count += 1
finally:
# 确保 Access 连接关闭
try:
@@ -230,7 +344,6 @@ def process_sync_task():
log_info(f"已关闭 Access 连接: {os.path.basename(clean_path)}")
except Exception as close_err:
log_warning(f"关闭 Access 连接时出错: {close_err}")
# 输出文件级别的汇总
if file_success_count > 0 or file_error_count > 0:
summary = f"文件 [{os.path.basename(clean_path)}] 同步汇总: "
@@ -245,16 +358,20 @@ def process_sync_task():
log_critical(f"全局异常: {e}")
return False
finally:
try: sql_conn.close()
except: pass
try:
sql_conn.close()
except:
pass
# ================= Uptime Kuma 心跳 =================
# 使用 uptime_kuma_utils.UptimeKumaMonitor 替代原有实现
if __name__ == "__main__":
log_start("增量同步服务已启动 (配置驱动模式)")
log_start("增量同步服务已启动 (配置驱动模式 + 逐主键落库校验)")
log_info(f"轮询间隔: {POLL_INTERVAL}")
log_info(f"监控配置: {len(SYNC_MAPPING)} 个文件")
log_info(f"提交后复核: {'开启' if ENABLE_POST_COMMIT_VERIFY else '关闭'}")
if UPTIME_KUMA_CONFIG.get('enabled', False):
log_info(f"心跳间隔: {UPTIME_KUMA_CONFIG['heartbeat_interval']}")
log_info("=" * 70)
@@ -271,7 +388,7 @@ if __name__ == "__main__":
uptime_monitor.check_and_send_heartbeat()
# 如果有工作,说明可能还有积压,休息短一点(0.1s)
# 如果没工作,休息标准间隔(5s)
# 如果没工作,休息标准间隔
time.sleep(0.1 if has_work else POLL_INTERVAL)
except KeyboardInterrupt: