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
相关产品推荐
相关产品推荐

