Add identity column check and manage IDENTITY_INSERT state during sync
This commit is contained in:
@@ -15,6 +15,20 @@ LoggerManager("run_incremental_sync", log_prefix="incremental")
|
|||||||
uptime_monitor = UptimeKumaMonitor(UPTIME_KUMA_CONFIG)
|
uptime_monitor = UptimeKumaMonitor(UPTIME_KUMA_CONFIG)
|
||||||
uptime_monitor.set_logger(log_warning)
|
uptime_monitor.set_logger(log_warning)
|
||||||
|
|
||||||
|
# ================= 辅助函数 =================
|
||||||
|
|
||||||
|
def has_identity_column(sql_cursor, schema, table):
|
||||||
|
"""检查表是否包含标识列"""
|
||||||
|
query = """
|
||||||
|
SELECT COUNT(*)
|
||||||
|
FROM sys.columns c
|
||||||
|
JOIN sys.tables t ON c.object_id = t.object_id
|
||||||
|
JOIN sys.schemas s ON t.schema_id = s.schema_id
|
||||||
|
WHERE s.name = ? AND t.name = ? AND c.is_identity = 1
|
||||||
|
"""
|
||||||
|
sql_cursor.execute(query, (schema, table))
|
||||||
|
return sql_cursor.fetchone()[0] > 0
|
||||||
|
|
||||||
# ================= 主逻辑 =================
|
# ================= 主逻辑 =================
|
||||||
|
|
||||||
def process_sync_task():
|
def process_sync_task():
|
||||||
@@ -141,13 +155,27 @@ def process_sync_task():
|
|||||||
# 2. 插入
|
# 2. 插入
|
||||||
if new_rows:
|
if new_rows:
|
||||||
insert_sql = db_utils.generate_insert_sql(target_schema, target_table, acc_cols)
|
insert_sql = db_utils.generate_insert_sql(target_schema, target_table, acc_cols)
|
||||||
try: sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} ON")
|
|
||||||
except: pass
|
# 检查并启用 IDENTITY_INSERT
|
||||||
|
has_identity = has_identity_column(sql_cursor, target_schema, target_table)
|
||||||
|
identity_enabled = False
|
||||||
|
|
||||||
|
if has_identity:
|
||||||
|
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)
|
sql_cursor.executemany(insert_sql, new_rows)
|
||||||
|
|
||||||
try: sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
|
# 关闭 IDENTITY_INSERT
|
||||||
except: pass
|
if identity_enabled:
|
||||||
|
try:
|
||||||
|
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
|
||||||
|
except Exception as id_err:
|
||||||
|
log_warning(f"关闭 IDENTITY_INSERT 时警告: {id_err}")
|
||||||
|
|
||||||
# 3. 标记日志 Synced = 1
|
# 3. 标记日志 Synced = 1
|
||||||
log_placeholders = ','.join(['?'] * len(log_ids))
|
log_placeholders = ','.join(['?'] * len(log_ids))
|
||||||
@@ -165,6 +193,11 @@ def process_sync_task():
|
|||||||
file_total_records += len(record_ids)
|
file_total_records += len(record_ids)
|
||||||
|
|
||||||
except Exception as tbl_err:
|
except Exception as tbl_err:
|
||||||
|
# 确保清理 IDENTITY_INSERT 状态
|
||||||
|
try:
|
||||||
|
sql_cursor.execute(f"SET IDENTITY_INSERT {target_full_name} OFF")
|
||||||
|
except:
|
||||||
|
pass
|
||||||
log_error(f"表 [{acc_table}] 同步失败: {tbl_err}")
|
log_error(f"表 [{acc_table}] 同步失败: {tbl_err}")
|
||||||
sql_conn.rollback()
|
sql_conn.rollback()
|
||||||
file_error_count += 1
|
file_error_count += 1
|
||||||
|
|||||||
Reference in New Issue
Block a user