You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SQLAlchemy ORM 2.0如何捕获批量更新的字段新旧值实现审计?

SQLAlchemy ORM 2.0 审计功能实现(含批量更新追踪)

在ORM 2.0中,实例级的增删改审计可以通过属性变更追踪实现,但批量SQL语句(如直接构造update/delete)不会加载实例,无法用get_history获取新旧值。下面分场景给出完整实现方案:

一、实例级增删改审计(基础实现)

对于先查询实例再修改的场景,利用Session的before_flush事件和get_history方法追踪字段变更,生成{列名: [旧值, 新值]}格式的审计字典:

from sqlalchemy import create_engine, Column, Integer, String
from sqlalchemy.orm import sessionmaker, declarative_base
from sqlalchemy.orm.attributes import get_history

Base = declarative_base()

class User(Base):
    __tablename__ = 'users'
    id = Column(Integer, primary_key=True)
    name = Column(String(50))
    email = Column(String(100))

# 初始化数据库连接
engine = create_engine('sqlite:///audit.db')
Base.metadata.create_all(engine)
Session = sessionmaker(bind=engine)

def audit_instance_operations(session, flush_context, instances):
    # 处理新增记录
    for instance in session.new:
        audit_data = {col.name: [None, getattr(instance, col.name)] 
                      for col in instance.__table__.columns}
        print(f"[新增审计] ID:{instance.id} 变更: {audit_data}")
    
    # 处理更新记录
    for instance in session.dirty:
        audit_data = {}
        for col in instance.__table__.columns:
            col_name = col.name
            history = get_history(instance, col_name)
            if history.has_changes():
                old_val = history.deleted[0] if history.deleted else None
                new_val = history.added[0] if history.added else None
                audit_data[col_name] = [old_val, new_val]
        if audit_data:
            print(f"[更新审计] ID:{instance.id} 变更: {audit_data}")
    
    # 处理删除记录
    for instance in session.deleted:
        audit_data = {col.name: [getattr(instance, col.name), None] 
                      for col in instance.__table__.columns}
        print(f"[删除审计] ID:{instance.id} 变更: {audit_data}")

# 绑定审计事件到Session
Session.configure(before_flush=audit_instance_operations)

# 测试实例修改
with Session() as session:
    # 新增
    new_user = User(name="Alice", email="alice@test.com")
    session.add(new_user)
    # 更新
    existing_user = session.query(User).filter_by(id=1).first()
    if existing_user:
        existing_user.name = "Alice Updated"
        existing_user.email = "alice_updated@test.com"
    # 删除
    delete_user = session.query(User).filter_by(id=2).first()
    if delete_user:
        session.delete(delete_user)
    session.commit()

二、批量UPDATE语句的审计实现

批量更新直接生成SQL执行,不会加载实例,核心思路是先查询待更新记录的旧值,再与更新的新值对比生成审计数据。提供两种实现方式:

方式1:封装业务工具函数(可控性强)

在业务代码中封装批量更新方法,强制先查询旧值再执行更新:

def bulk_update_with_audit(session, model, update_values, filter_condition):
    # 1. 查询待更新记录的旧值
    target_records = session.query(model).filter(filter_condition).all()
    if not target_records:
        return 0
    
    # 2. 执行批量更新(synchronize_session=False避免自动加载实例)
    update_count = session.query(model).filter(filter_condition).update(
        update_values,
        synchronize_session=False
    )
    
    # 3. 生成审计数据
    audit_records = []
    for record in target_records:
        change_data = {}
        for col_name, new_val in update_values.items():
            old_val = getattr(record, col_name)
            if old_val != new_val:
                change_data[col_name] = [old_val, new_val]
        if change_data:
            audit_records.append({
                "record_id": record.id,
                "changes": change_data
            })
    
    # 输出/存储审计记录
    if audit_records:
        print(f"[批量更新审计] 共{len(audit_records)}条记录变更: {audit_records}")
    return update_count

# 测试批量更新
with Session() as session:
    bulk_update_with_audit(
        session=session,
        model=User,
        update_values={"name": "Bulk Updated", "email": "bulk@test.com"},
        filter_condition=User.id.in_([1, 3])
    )
    session.commit()

方式2:全局事件拦截(无侵入业务代码)

利用SQLAlchemy的before_execute事件拦截所有UPDATE语句,自动解析并生成审计数据:

from sqlalchemy import event
from sqlalchemy.sql import Update

def intercept_bulk_updates(conn, clauseelement, multiparams, params):
    # 仅拦截UPDATE语句
    if not isinstance(clauseelement, Update):
        return
    
    target_table = clauseelement.table
    filter_clause = clauseelement.whereclause
    if not filter_clause:
        # 避免全表更新的风险,可根据需求调整
        print("[警告] 拦截到无过滤条件的全表更新,跳过审计")
        return
    
    # 解析更新的字段与新值
    update_values = {}
    for col, value_expr in clauseelement.values.items():
        # 处理硬编码值或绑定参数
        if hasattr(value_expr, 'value'):
            update_values[col.name] = value_expr.value
        else:
            # 从参数中提取实际值
            update_values[col.name] = multiparams[0][col.name]
    
    # 创建临时Session查询旧值
    temp_session = Session(bind=conn)
    target_records = temp_session.query(target_table).filter(filter_clause).all()
    temp_session.close()
    
    # 生成审计记录
    audit_records = []
    for record in target_records:
        change_data = {}
        for col_name, new_val in update_values.items():
            old_val = getattr(record, col_name)
            if old_val != new_val:
                change_data[col_name] = [old_val, new_val]
        if change_data:
            audit_records.append({
                "record_id": record.id,
                "changes": change_data
            })
    
    if audit_records:
        print(f"[拦截式批量更新审计] 共{len(audit_records)}条记录变更: {audit_records}")

# 绑定事件到数据库引擎
event.listen(engine, "before_execute", intercept_bulk_updates)

# 测试直接执行批量UPDATE
with Session() as session:
    session.query(User).filter(User.id.in_([1,3])).update(
        {"name": "Intercepted Bulk", "email": "intercepted@test.com"},
        synchronize_session=False
    )
    session.commit()

三、批量DELETE语句的审计

类似批量更新,先查询待删除记录的旧值,再执行删除并生成审计数据:

def bulk_delete_with_audit(session, model, filter_condition):
    target_records = session.query(model).filter(filter_condition).all()
    if not target_records:
        return 0
    
    delete_count = session.query(model).filter(filter_condition).delete(
        synchronize_session=False
    )
    
    audit_records = []
    for record in target_records:
        change_data = {col.name: [getattr(record, col.name), None] 
                      for col in record.__table__.columns}
        audit_records.append({
            "record_id": record.id,
            "changes": change_data
        })
    
    if audit_records:
        print(f"[批量删除审计] 共{len(audit_records)}条记录被删除: {audit_records}")
    return delete_count

# 测试批量删除
with Session() as session:
    bulk_delete_with_audit(
        session=session,
        model=User,
        filter_condition=User.id.in_([2,4])
    )
    session.commit()

内容的提问来源于stack exchange,提问作者Manika Midha

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 15:54:51