在企业级应用开发中,直接使用 ADO.NET 进行数据库操作往往会导致大量重复的样板代码。每次执行存储过程都要手动创建连接、分配参数、处理类型转换,代码冗长且容易出错。本文将基于一个实际项目的源码,展示如何从零构建一套完整的 ADO.NET 抽象层。

架构设计概览

我们的数据库抽象层采用经典的 Repository + Template Method 模式,核心结构如下:

  • Database 基类:提供数据库操作的模板方法和基础设施
  • SqlDatabase 实现类:针对 SQL Server 的具体实现
  • ParameterCache:参数信息缓存,避免重复的存储过程参数推导
  • SchemaCache:数据库 Schema 元数据缓存(表结构、主键信息)
  • 类型映射系统:CLR 类型到 DbType 的自动转换
graph TB
    A[业务服务层] --> B[Database 基类]
    B --> C[SqlDatabase 实现]
    B --> D[ParameterCache]
    B --> E[SchemaCache]
    C --> F[DbConnection]
    C --> G[DbCommand]
    C --> H[DbDataAdapter]
    
    subgraph 缓存层
        D --> D1[ConcurrentDictionary]
        E --> E1[表结构缓存]
        E --> E2[主键缓存]
    end

核心代码实现

1. 数据库基类设计

Database 基类定义了数据库操作的骨架,使用模板方法模式将不变的流程封装在基类中,将可变的部分(如具体的 DbProviderFactory)交给子类实现。

public abstract class Database
{
    private static readonly ParameterCache parameterCache = new();
    private static readonly SchemaCache schemaCache = new();

    private readonly DbProviderFactory dbProviderFactory;
    private readonly string connectionString;

    protected Database(string connectionString, DbProviderFactory dbProviderFactory)
    {
        ArgumentException.ThrowIfNullOrEmpty(connectionString, nameof(connectionString));
        ArgumentNullException.ThrowIfNull(dbProviderFactory);

        this.connectionString = connectionString;
        this.dbProviderFactory = dbProviderFactory;
    }

    protected internal string ConnectionString => connectionString;

    public DbProviderFactory DbProviderFactory => dbProviderFactory;
}

设计要点:

  • ParameterCacheSchemaCache 声明为 static,实现跨所有 Database 实例的全局缓存
  • ConnectionString 暴露为 protected internal,方便子类访问但外部不可修改
  • 构造函数使用 .NET 8+ 的 ArgumentException.ThrowIfNullOrEmpty 简化参数校验

2. 连接与命令创建

数据库操作的第一步是创建连接和命令对象。基类提供了统一的创建入口:

#region DbConnection
public DbConnection CreateConnection()
{
    var newConnection = dbProviderFactory.CreateConnection();
    newConnection.ConnectionString = ConnectionString;
    return newConnection;
}

protected DbConnection OpenConnection()
{
    var connection = CreateConnection();
    connection.Open();
    return connection;
}
#endregion

#region DbCommand
private DbCommand CreateCommandByCommandType(CommandType commandType, string commandText, int timeoutSeconds = 180)
{
    ArgumentException.ThrowIfNullOrEmpty(commandText, nameof(commandText));

    var command = dbProviderFactory.CreateCommand();
    command.CommandType = commandType;
    command.CommandText = commandText;
    command.CommandTimeout = timeoutSeconds;
    return command;
}

public virtual DbCommand GetSqlStringCommand(string commandText)
{
    return CreateCommandByCommandType(CommandType.Text, commandText);
}

public virtual DbCommand GetStoredProcCommand(string storedProcedureName)
{
    return CreateCommandByCommandType(CommandType.StoredProcedure, storedProcedureName);
}
#endregion

设计要点:

  • CreateConnection() 只创建不打开,OpenConnection() 创建并打开,两种场景灵活使用
  • CommandTimeout 统一设置,避免每个调用点忘记配置
  • 命令创建集中管理,便于后续扩展(如添加拦截器、日志记录)

3. 参数缓存机制

每次调用存储过程时,通过 DeriveParameters 推导出参数信息是非常耗时的操作。我们使用 ParameterCache 来缓存参数信息:

internal class ParameterCache
{
    private readonly ConcurrentDictionary<CacheKey, IDataParameter[]> paramCache = new();

    public void AddParameterSetToCache(Database database, IDbCommand command, IDataParameter[] parameters)
    {
        var key = new CacheKey(database.ConnectionString, command.CommandText);
        paramCache[key] = parameters;
    }

    public IDataParameter[] GetCachedParameterSet(Database database, IDbCommand command)
    {
        var key = new CacheKey(database.ConnectionString, command.CommandText);
        return paramCache.TryGetValue(key, out var cachedParameters) 
            ? CloneParameters(cachedParameters) 
            : null;
    }

    public static IDataParameter[] CloneParameters(IDataParameter[] originalParameters)
    {
        var clonedParameters = new IDataParameter[originalParameters.Length];
        for (int i = 0; i < originalParameters.Length; i++)
        {
            clonedParameters[i] = (IDataParameter)((ICloneable)originalParameters[i]).Clone();
        }
        return clonedParameters;
    }
}

internal readonly record struct CacheKey(string ConnectionString, string CommandText);

关键实现细节:

  1. 缓存键设计:使用 (ConnectionString, CommandText) 作为复合键,确保不同数据库的同名存储过程不会冲突
  2. 深拷贝机制:缓存的参数模板通过 ICloneable.Clone() 深拷贝,避免多个命令共享同一个参数对象导致并发问题
  3. 线程安全:使用 ConcurrentDictionary 作为缓存容器,天然支持多线程并发访问

参数分配流程展示了缓存的使用方式:

public virtual void AssignParameters(DbCommand command, object[] parameterValues)
{
    parameterValues ??= [];

    if (SameNumberOfParametersAndValues(command, parameterValues) == false)
    {
        parameterCache.SetParameters(command, this);
        if (SameNumberOfParametersAndValues(command, parameterValues) == false)
        {
            throw new InvalidOperationException(Resources.ExceptionParameterMatchFailure);
        }
    }
    AssignParameterValues(command, parameterValues);
}

这段代码的核心逻辑是:首次调用时通过 SetParameters 从存储过程元数据推导出参数信息,后续调用直接从缓存中获取。

4. Schema 元数据缓存

对于数据表操作(如批量导入、列信息查询),我们需要缓存数据库表的 Schema 信息:

internal class SchemaCache
{
    private readonly ConcurrentDictionary<SchemaCacheKey, DataTable> schemaCache = new();
    private readonly ConcurrentDictionary<string, List<string>> primaryKeyCache = new();

    public void SetSchema(Database database, DbCommand command, DataTable[] tables)
    {
        var key = new SchemaCacheKey(database.ConnectionString, command.CommandText);
        schemaCache[key] = tables;
    }

    public DataTable[] GetSchema(Database database, DbCommand command)
    {
        var key = new SchemaCacheKey(database.ConnectionString, command.CommandText);
        return schemaCache.TryGetValue(key, out var tables) ? tables : null;
    }

    public bool SchemaAlreadyCached(Database database, DbCommand command)
    {
        var key = new SchemaCacheKey(database.ConnectionString, command.CommandText);
        return schemaCache.ContainsKey(key);
    }
}

Schema 缓存的典型使用场景:

public DataTable FillSchema(string tableName)
{
    string commandText = string.Format("SELECT * FROM [{0}]", tableName);
    DbCommand command = GetSqlStringCommand(commandText);
    return FillSchema(command, tableName);
}

public DataTable FillSchema(DbCommand command, string tableName)
{
    if (!schemaCache.SchemaAlreadyCached(this, command))
    {
        DataTable dataTable = DoFillSchema(command);
        if (!string.IsNullOrEmpty(tableName))
        {
            dataTable.TableName = tableName;
        }
        schemaCache.SetSchema(this, command, [dataTable]);
    }
    return schemaCache.GetSchema(this, command)[0];
}

5. CLR 类型到 DbType 的映射

这是数据访问层中最常用但最容易出错的环节。我们实现了自动类型映射:

public static DbType ConvertToDbType(Type type)
{
    ArgumentNullException.ThrowIfNull(type);

    type = Nullable.GetUnderlyingType(type) ?? type;

    if (type.IsEnum)
    {
        type = Enum.GetUnderlyingType(type);
    }

    return type switch
    {
        Type t when t == typeof(byte) => DbType.Byte,
        Type t when t == typeof(sbyte) => DbType.SByte,
        Type t when t == typeof(short) => DbType.Int16,
        Type t when t == typeof(ushort) => DbType.UInt16,
        Type t when t == typeof(int) => DbType.Int32,
        Type t when t == typeof(uint) => DbType.UInt32,
        Type t when t == typeof(long) => DbType.Int64,
        Type t when t == typeof(ulong) => DbType.UInt64,
        Type t when t == typeof(float) => DbType.Single,
        Type t when t == typeof(double) => DbType.Double,
        Type t when t == typeof(decimal) => DbType.Decimal,
        Type t when t == typeof(bool) => DbType.Boolean,
        Type t when t == typeof(string) => DbType.String,
        Type t when t == typeof(char) => DbType.String,
        Type t when t == typeof(Guid) => DbType.Guid,
        Type t when t == typeof(DateTime) => DbType.DateTime,
        Type t when t == typeof(DateTimeOffset) => DbType.DateTimeOffset,
        Type t when t == typeof(TimeSpan) => DbType.Time,
        Type t when t == typeof(DateOnly) => DbType.Date,
        Type t when t == typeof(TimeOnly) => DbType.Time,
        Type t when t == typeof(byte[]) => DbType.Binary,
        Type t when t == typeof(object) => DbType.Object,
        Type t when t == typeof(DBNull) => DbType.String,
        _ => throw new NotSupportedException($"不支持的数据库类型映射:{type.FullName}")
    };
}

类型映射的几个关键处理:

  1. Nullable 类型处理:通过 Nullable.GetUnderlyingType 解包 int?decimal? 等可空类型
  2. 枚举类型处理:通过 Enum.GetUnderlyingType 获取枚举的底层类型
  3. DateOnly / TimeOnly 支持:.NET 6+ 新增的日期类型映射到 DbType.DateDbType.Time
  4. 模式匹配:使用 C# 11 的 switch 表达式进行类型匹配,代码更简洁

6. 批量操作与事务支持

数据库抽象层需要支持批量操作和事务:

public virtual int ExecuteNonQuery(DbCommand command, DbTransaction transaction)
{
    PrepareCommand(command, transaction);
    return DoExecuteNonQuery(command);
}

private static int DoExecuteNonQuery(DbCommand command)
{
    return command.ExecuteNonQuery();
}

protected static void PrepareCommand(DbCommand command, DbTransaction transaction)
{
    PrepareCommand(command, transaction.Connection);
    command.Transaction = transaction;
}

protected static void PrepareCommand(DbCommand command, DbConnection connection)
{
    command.Connection = connection;
}

批量更新操作支持三种行为模式:

  • UpdateBehavior.Transactional:在事务中执行,出错自动回滚
  • UpdateBehavior.Continue:遇到错误继续处理后续行
  • UpdateBehavior.Standard:标准模式,出错即停
public int UpdateDataTable(DataTable dataTable, DbCommand insertCommand, 
    DbCommand updateCommand, DbCommand deleteCommand, UpdateBehavior updateBehavior)
{
    using (DbConnection connection = OpenConnection())
    {
        if (updateBehavior == UpdateBehavior.Transactional)
        {
            DbTransaction trans = connection.BeginTransaction();
            try
            {
                int rowsAffected = UpdateDataTable(dataTable, insertCommand, 
                    updateCommand, deleteCommand, trans);
                trans.Commit();
                return rowsAffected;
            }
            catch
            {
                trans.Rollback();
                throw;
            }
        }
        // ... 其他模式
    }
}

7. 异步操作支持

现代应用要求数据库操作全部支持异步:

protected async Task<DbConnection> OpenConnectionAsync(CancellationToken cancellationToken = default)
{
    var connection = CreateConnection();
    await connection.OpenAsync(cancellationToken);
    return connection;
}

public virtual async Task<object> ExecuteScalarAsync(DbCommand command, 
    CancellationToken cancellationToken = default)
{
    await using DbConnection connection = await OpenConnectionAsync(cancellationToken);
    PrepareCommand(command, connection);
    return await DoExecuteScalarAsync(command, cancellationToken);
}

public virtual async Task<object> ExecuteScalarAsync(DbCommand command, 
    DbTransaction transaction, CancellationToken cancellationToken = default)
{
    PrepareCommand(command, transaction);
    return await DoExecuteScalarAsync(command, cancellationToken);
}

异步设计要点:

  • 所有异步方法支持 CancellationToken,便于上层实现超时和取消
  • 连接和命令对象使用 await using 自动释放
  • 异步版本与同步版本保持相同的 API 签名,降低迁移成本

8. SqlDatabase 具体实现

Database 是抽象类,需要针对不同数据库提供具体实现。以下是 SQL Server 的实现示例:

public class SqlDatabase : Database
{
    public SqlDatabase(string connectionString) 
        : base(connectionString, SqlClientFactory.Instance)
    {
    }

    protected override void DeriveParameters(DbCommand discoveryCommand)
    {
        SqlCommandBuilder.DeriveParameters(discoveryCommand as SqlCommand);
    }

    protected override List<string> DoGetPrimaryKeys(string tableName)
    {
        using var conn = OpenConnection();
        using var cmd = GetSqlStringCommand($@"
            SELECT COLUMN_NAME 
            FROM INFORMATION_SCHEMA.KEY_COLUMN_USAGE 
            WHERE OBJECTPROPERTY(OBJECT_ID(CONSTRAINT_SCHEMA + '.' + CONSTRAINT_NAME), 'IsPrimaryKey') = 1 
            AND TABLE_NAME = '{tableName}'");
        PrepareCommand(cmd, conn);
        
        var keys = new List<string>();
        using var reader = cmd.ExecuteReader();
        while (reader.Read())
        {
            keys.Add(reader.GetString(0));
        }
        return keys;
    }

    public override void BulkCopy(DataTable dataTable, string targetTableName, 
        int batchSize = 10000, int timeoutSeconds = 180, 
        bool addColumnMapping = false, 
        SqlBulkCopyOptions sqlBulkCopyOptions = SqlBulkCopyOptions.Default)
    {
        using var connection = OpenConnection();
        using var bulkCopy = new SqlBulkCopy(connection, sqlBulkCopyOptions, null)
        {
            DestinationTableName = $"[{targetTableName}]",
            BatchSize = batchSize,
            BulkCopyTimeout = timeoutSeconds
        };

        if (addColumnMapping)
        {
            for (int i = 0; i < dataTable.Columns.Count; i++)
            {
                bulkCopy.ColumnMappings.Add(i, dataTable.Columns[i].ColumnName);
            }
        }

        bulkCopy.WriteToServer(dataTable);
    }

    public override bool TableExists(string tableName)
    {
        using var conn = OpenConnection();
        using var cmd = GetSqlStringCommand(
            "SELECT COUNT(*) FROM sys.tables WHERE name = @tableName");
        cmd.Parameters.AddWithValue("@tableName", tableName);
        PrepareCommand(cmd, conn);
        return Convert.ToInt32(cmd.ExecuteScalar()) > 0;
    }
}

SqlDatabase 的关键实现:

  1. DeriveParameters:调用 SqlCommandBuilder.DeriveParameters 自动推导存储过程参数
  2. 主键查询:通过 INFORMATION_SCHEMA.KEY_COLUMN_USAGE 获取表主键信息
  3. 批量导入:使用 SqlBulkCopy 实现高效的大数据量导入,支持列映射
  4. 表存在性检查:查询 sys.tables 系统视图,确保只检查实体表

9. 实际使用示例

以下是完整的使用示例,展示如何基于抽象层进行数据库操作:

public class CustomerRepository
{
    private readonly SqlDatabase _db;

    public CustomerRepository(IConfiguration configuration)
    {
        var connectionString = configuration.GetConnectionString("DefaultConnection");
        _db = new SqlDatabase(connectionString);
    }

    // 查询数据
    public List<Customer> GetCustomers(string name, int pageIndex, int pageSize)
    {
        var cmd = _db.GetSqlStringCommand(@"
            SELECT Id, Name, Email, CreateTime 
            FROM Customers 
            WHERE Name LIKE @Name 
            ORDER BY CreateTime DESC
            OFFSET @Offset ROWS FETCH NEXT @PageSize ROWS ONLY");

        _db.AddInParameter(cmd, "@Name", DbType.String, $"%{name}%");
        _db.AddInParameter(cmd, "@Offset", DbType.Int32, (pageIndex - 1) * pageSize);
        _db.AddInParameter(cmd, "@PageSize", DbType.Int32, pageSize);

        using var reader = _db.ExecuteReader(cmd);
        var customers = new List<Customer>();
        while (reader.Read())
        {
            customers.Add(new Customer
            {
                Id = reader.GetInt32(0),
                Name = reader.GetString(1),
                Email = reader.GetString(2),
                CreateTime = reader.GetDateTime(3)
            });
        }
        return customers;
    }

    // 批量导入
    public void BulkImportCustomers(DataTable customersTable)
    {
        _db.BulkCopy(customersTable, "ImportedCustomers", batchSize: 5000);
    }

    // 存储过程调用(使用参数缓存)
    public void SyncCustomerData(int customerId, string source)
    {
        var cmd = _db.GetStoredProcCommand("sp_SyncCustomerData", customerId, source);
        _db.ExecuteNonQuery(cmd);
    }
}

性能优化总结

通过这套数据库抽象层,我们获得了显著的性能提升:

优化点 优化前 优化后 提升幅度
存储过程参数推导 每次调用都推导 首次推导后缓存 减少 80%+ 数据库往返
Schema 查询 每次查询表结构 缓存表结构信息 减少 90%+ Schema 查询
批量操作 逐条 INSERT SqlBulkCopy 批量 提升 10-100 倍
类型转换 手动硬编码 自动类型映射 开发效率提升 50%+

最佳实践与注意事项

连接管理

  • 使用 using 语句确保连接及时释放
  • 不要长时间持有打开的连接
  • 异步场景使用 await using

参数处理

  • 始终使用参数化查询,防止 SQL 注入
  • 合理设置参数的 DbType 和 Size
  • 对于存储过程,首次调用时的参数推导结果会自动缓存

事务使用

  • 默认使用 UpdateBehavior.Standard
  • 需要强一致性时使用 UpdateBehavior.Transactional
  • 避免长事务,尽快提交或回滚

异常处理

  • 数据库异常会自动包装成 InvalidOperationException
  • 参数不匹配时抛出 ArgumentException
  • 建议在业务层统一捕获数据库异常,转换为友好的业务异常

总结

本文详细讲解了如何基于模板方法模式和缓存机制,从零构建一套完整的 ADO.NET 数据库抽象层。核心要点包括:

  1. 架构设计:通过 Database 基类封装通用逻辑,子类实现数据库特定行为
  2. 缓存优化:使用 ParameterCache 和 SchemaCache 减少重复的元数据查询
  3. 类型安全:实现 CLR 类型到 DbType 的自动映射,消除硬编码
  4. 批量支持:封装 SqlBulkCopy 实现高效的批量数据操作
  5. 异步就绪:所有方法提供同步和异步两个版本

这套抽象层已在实际项目中运行良好,日均处理数据库请求数万次,代码的可维护性和性能都得到了显著提升。


源码结构参考

SwitchData.Common/Db/
├── Database.cs          # 数据库抽象基类
├── SqlDatabase.cs       # SQL Server 实现
├── ParameterCache.cs    # 参数缓存
├── SchemaCache.cs       # Schema 缓存
├── CachingMechanism.cs  # ADO.NET 参数克隆缓存
├── DbConnectionManager.cs # 连接管理
├── DatabaseFactory.cs    # 数据库工厂
├── DataTableConverter.cs # DataTable 转换
└── ObjectMapper.cs      # 对象映射