EF Core中通过拦截器动态设置search_path的异常排查
背景与问题描述
在多租户场景中,通过动态修改连接字符串设置search_path会导致每个租户生成独立的数据库连接池,造成资源浪费。因此改用EF Core拦截器,为每个执行的SQL语句前置SET search_path TO "some_schema";语句,从DbContext动态获取租户对应的schema名称。
使用环境:.NET 6 + EF Core + PostgreSQL
实现的拦截器代码如下:
public class SchemaInterceptor : DbCommandInterceptor { private readonly string path; public SchemaInterceptor(string path) { this.path = path; } private void SetSearchPath(DbCommand command) { command.CommandText = $"SET search_path TO \"{path}\";\n{command.CommandText}"; } public override InterceptionResult<DbDataReader> ReaderExecuting(DbCommand command, CommandEventData eventData, InterceptionResult<DbDataReader> result) { SetSearchPath(command); return base.ReaderExecuting(command, eventData, result); } public override ValueTask<InterceptionResult<DbDataReader>> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult<DbDataReader> result, CancellationToken cancellationToken = default) { SetSearchPath(command); return base.ReaderExecutingAsync(command, eventData, result, cancellationToken); } public override InterceptionResult<int> NonQueryExecuting(DbCommand command, CommandEventData eventData, InterceptionResult<int> result) { SetSearchPath(command); return base.NonQueryExecuting(command, eventData, result); } public override ValueTask<InterceptionResult<int>> NonQueryExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult<int> result, CancellationToken cancellationToken = default) { SetSearchPath(command); return base.NonQueryExecutingAsync(command, eventData, result, cancellationToken); } public override InterceptionResult<object> ScalarExecuting(DbCommand command, CommandEventData eventData, InterceptionResult<object> result) { SetSearchPath(command); return base.ScalarExecuting(command, eventData, result); } public override ValueTask<InterceptionResult<object>> ScalarExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult<object> result, CancellationToken cancellationToken = default) { SetSearchPath(command); return base.ScalarExecutingAsync(command, eventData, result, cancellationToken); } }
该拦截器对SELECT和简单INSERT操作有效,但插入带有多对多关联的实体时失败,报错信息如下:
[15:21:03 DBG] Executing DbCommand [Parameters=[@p0='?' (DbType = DateTime), @p1='?' (DbType = DateTime), @p2='?', @p3='?'], CommandType='Text', CommandTimeout='30'] INSERT INTO projects (date_added, date_last_modified, name, notes) VALUES (@p0, @p1, @p2, @p3) RETURNING projects_id; [15:21:03 INF] Executed DbCommand (1ms) [Parameters=[@p0='?' (DbType = DateTime), @p1='?' (DbType = DateTime), @p2='?', @p3='?'], CommandType='Text', CommandTimeout='30'] SET search_path TO "test"; INSERT INTO projects (date_added, date_last_modified, name, notes) VALUES (@p0, @p1, @p2, @p3) RETURNING projects_id; [15:21:03 DBG] The foreign key property 'Project.Id' was detected as changed. Consider using 'DbContextOptionsBuilder.EnableSensitiveDataLogging' to see property values. 3 [15:21:03 DBG] The foreign key property 'projects_tags.projects_id' was detected as changed. Consider using 'DbContextOptionsBuilder.EnableSensitiveDataLogging' to see property values. [15:21:03 DBG] A data reader was disposed. [15:21:03 DBG] Executing 3 update commands as a batch. [15:21:03 DBG] Creating DbCommand for 'ExecuteReader'. [15:21:03 DBG] Created DbCommand for 'ExecuteReader' (0ms). INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p4, @p5); INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p6, @p7); INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p8, @p9); [15:21:03 INF] Executed DbCommand (8ms) [Parameters=[@p4='?' (DbType = Int32), @p5='?' (DbType = Int32), @p6='?' (DbType = Int32), @p7='?' (DbType = Int32), @p8='?' (DbType = Int32), @p9='?' (DbType = Int32)], CommandType='Text', CommandTimeout='30'] SET search_path TO "test"; INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p4, @p5); INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p6, @p7); INSERT INTO projects_tags (project_tags_id, projects_id) VALUES (@p8, @p9); [...] [15:21:03 DBG] Microsoft.EntityFrameworkCore.DbUpdateConcurrencyException: The database operation was expected to affect 1 row(s), but actually affected 0 row(s); data may have been modified or deleted since entities were loaded. at Npgsql.EntityFrameworkCore.PostgreSQL.Update.Internal.NpgsqlModificationCommandBatch.ConsumeAsync(RelationalDataReader reader, CancellationToken cancellationToken) [...]
实际未启用并发冲突检测,且直接在DbContext中设置search_path时代码可正常运行,仅插入无关联实体时拦截器能正常工作。
问题原因分析
EF Core处理批量更新/插入(如多对多关联的中间表插入)时,会期望数据库返回的受影响行数与命令数量严格匹配。原拦截器在批量命令前添加SET search_path语句后,该语句本身会返回0行受影响的结果,导致EF Core的批量命令处理器接收到的结果集结构和预期不符,从而误判为并发冲突。
从日志可见,批量插入3条中间表记录时,实际执行的SQL先执行SET语句,再执行3条INSERT,数据库返回的结果集中第一条是SET的0行影响,后面才是3条INSERT的各1行影响,但EF Core期望直接返回3条结果,因此触发“预期影响1行实际0行”的错误。
修复方案
利用PostgreSQL的会话级特性——SET search_path的效果在整个数据库连接会话中生效,无需为每条SQL重复添加。改用DbConnectionInterceptor在连接打开时执行一次schema设置即可:
public class SchemaConnectionInterceptor : DbConnectionInterceptor { private readonly string _schemaName; public SchemaConnectionInterceptor(string schemaName) { _schemaName = schemaName; } public override async ValueTask<InterceptionResult> ConnectionOpeningAsync(DbConnection connection, ConnectionEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) { if (connection.State != ConnectionState.Open) { await connection.OpenAsync(cancellationToken); using var command = connection.CreateCommand(); command.CommandText = $"SET search_path TO \"{_schemaName}\";"; await command.ExecuteNonQueryAsync(cancellationToken); return InterceptionResult.Suppress; } return await base.ConnectionOpeningAsync(connection, eventData, result, cancellationToken); } public override InterceptionResult ConnectionOpening(DbConnection connection, ConnectionEventData eventData, InterceptionResult result) { if (connection.State != ConnectionState.Open) { connection.Open(); using var command = connection.CreateCommand(); command.CommandText = $"SET search_path TO \"{_schemaName}\";"; command.ExecuteNonQuery(); return InterceptionResult.Suppress; } return base.ConnectionOpening(connection, eventData, result); } }
注册DbContext时替换原拦截器:
services.AddDbContext<MyDbContext>((sp, options) => { var connectionString = sp.GetRequiredService<IConfiguration>().GetConnectionString("Default"); options.UseNpgsql(connectionString) .AddInterceptors(new SchemaConnectionInterceptor(GetCurrentTenantSchema())); });
方案优势
- 避免SQL冗余,提升执行性能:无需为每条SQL重复添加SET语句
- 适配EF Core批量逻辑:会话级schema设置不会干扰EF Core对命令结果的判断
- 资源利用更高效:复用连接池,不会因动态schema生成独立连接池
内容的提问来源于stack exchange,提问作者Fabian

