ETL数据整合层设计模式选型:ASP.NET Core多数据源数据收集系统实现咨询
适合你的多数据源采集系统的设计模式组合:仓库 + 管道
针对你这个基于ASP.NET Core的多数据源数据采集、持久化与转换系统,我强烈推荐组合使用仓库模式(Repository) 和 管道模式(Pipeline)。这两个模式互补性极强,刚好能覆盖你的核心需求:
- 仓库模式负责封装不同数据源的读写逻辑,让上层业务代码和具体数据源彻底解耦,后续新增数据源只需要加新的实现类就行。
- 管道模式把数据采集、验证、转换、持久化的流程拆成独立的可复用步骤,方便你灵活调整流程,也更容易测试单个环节。
下面我给你详细拆解每个模式的用法,附上伪代码,还有处理多数据源实例的DI方案。
1. 仓库模式:封装多数据源操作
首先,我们为不同类型的数据源定义抽象仓库接口,然后针对每个具体数据源实现对应的仓库类。这样上层代码只依赖抽象接口,不用关心底层是API、数据库还是CSV文件。
伪代码实现:
第一步:定义抽象仓库接口
// 产品属性数据的抽象仓库 public interface IProductAttributeRepository { Task<IEnumerable<ProductAttributeEntity>> FetchAsync(); Task PersistAsync(IEnumerable<ProductAttributeEntity> entities); } // 资产数据的抽象仓库 public interface IAssetRepository { Task<IEnumerable<AssetEntity>> FetchAsync(); Task PersistAsync(IEnumerable<AssetEntity> entities); } // 通用仓库接口(可选,用来统一基础操作) public interface IRepository<TEntity> { Task<IEnumerable<TEntity>> FetchAsync(); Task PersistAsync(IEnumerable<TEntity> entities); }
第二步:实现具体数据源的仓库类
比如针对DataSource1(假设是外部API):
public class DataSource1ProductAttributeRepository : IProductAttributeRepository { private readonly HttpClient _httpClient; public DataSource1ProductAttributeRepository(HttpClient httpClient) { _httpClient = httpClient; } public async Task<IEnumerable<ProductAttributeEntity>> FetchAsync() { // 调用外部API获取数据,映射到预定义的实体 var apiResponse = await _httpClient.GetFromJsonAsync<DataSource1ProductDto[]>("api/products/attributes"); return apiResponse.Select(dto => new ProductAttributeEntity { Id = dto.ProductId, Name = dto.ProductName, Category = dto.ProductCategory // 其他属性映射逻辑 }); } public async Task PersistAsync(IEnumerable<ProductAttributeEntity> entities) { // 将数据持久化到你的统一存储(比如本地SQL Server) using var dbContext = new AppDbContext(); dbContext.ProductAttributes.AddRange(entities); await dbContext.SaveChangesAsync(); } }
再比如DataSource2(假设是本地数据库的外部数据表):
public class DataSource2ProductAttributeRepository : IProductAttributeRepository { private readonly AppDbContext _dbContext; public DataSource2ProductAttributeRepository(AppDbContext dbContext) { _dbContext = dbContext; } public async Task<IEnumerable<ProductAttributeEntity>> FetchAsync() { // 从外部数据表读取数据 return await _dbContext.ExternalProductAttributes .Select(ext => new ProductAttributeEntity { Id = ext.ProductId, Name = ext.ProductName, Category = ext.Category }) .ToListAsync(); } public async Task PersistAsync(IEnumerable<ProductAttributeEntity> entities) { // 同样持久化到统一存储 _dbContext.ProductAttributes.AddRange(entities); await _dbContext.SaveChangesAsync(); } }
DataSource3的资产仓库实现类似,这里就不重复写了。
2. 管道模式:流程化处理数据
你的数据流程是「采集→验证→转换→持久化」,管道模式刚好能把这些步骤拆成独立的处理器,每个处理器只做一件事,方便扩展和维护。
伪代码实现:
第一步:定义管道基础接口和管道类
// 单个管道步骤的接口 public interface IPipelineStep<TInput, TOutput> { Task<TOutput> ProcessAsync(TInput input); } // 管道类,用来串联多个步骤 public class Pipeline<TInput, TOutput> { private readonly IEnumerable<IPipelineStep<object, object>> _steps; public Pipeline(IEnumerable<IPipelineStep<object, object>> steps) { _steps = steps; } public async Task<TOutput> ExecuteAsync(TInput input) { var currentInput = (object)input; foreach (var step in _steps) { currentInput = await step.ProcessAsync(currentInput); } return (TOutput)currentInput; } }
第二步:实现具体的管道步骤
// 数据采集步骤:从指定仓库拉取数据 public class DataFetchStep<TEntity> : IPipelineStep<IRepository<TEntity>, IEnumerable<TEntity>> { public async Task<IEnumerable<TEntity>> ProcessAsync(IRepository<TEntity> repository) { return await repository.FetchAsync(); } } // 数据验证步骤:过滤无效数据 public class DataValidationStep<TEntity> : IPipelineStep<IEnumerable<TEntity>, IEnumerable<TEntity>> { public async Task<IEnumerable<TEntity>> ProcessAsync(IEnumerable<TEntity> entities) { var validEntities = entities.Where(e => e.Id != Guid.Empty).ToList(); if (validEntities.Count != entities.Count()) { // 这里可以加日志记录,比如"过滤了X条无效数据" } return await Task.FromResult(validEntities); } } // 数据转换步骤:映射到统一抽象实体 public class ProductToUnifiedEntityStep : IPipelineStep<IEnumerable<ProductAttributeEntity>, IEnumerable<UnifiedProductEntity>> { public async Task<IEnumerable<UnifiedProductEntity>> ProcessAsync(IEnumerable<ProductAttributeEntity> entities) { return entities.Select(e => new UnifiedProductEntity { ProductId = e.Id, ProductName = e.Name, Category = e.Category, // 其他统一属性的映射 }); } } // 持久化步骤:把统一实体保存到存储 public class UnifiedDataPersistStep : IPipelineStep<IEnumerable<UnifiedProductEntity>, bool> { private readonly AppDbContext _dbContext; public UnifiedDataPersistStep(AppDbContext dbContext) { _dbContext = dbContext; } public async Task<bool> ProcessAsync(IEnumerable<UnifiedProductEntity> entities) { _dbContext.UnifiedProducts.AddRange(entities); await _dbContext.SaveChangesAsync(); return true; } }
第三步:在业务服务中使用管道
public class ProductDataProcessingService { private readonly Pipeline<IRepository<ProductAttributeEntity>, bool> _productPipeline; public ProductDataProcessingService(Pipeline<IRepository<ProductAttributeEntity>, bool> productPipeline) { _productPipeline = productPipeline; } public async Task ProcessProductDataAsync(IRepository<ProductAttributeEntity> repository) { var success = await _productPipeline.ExecuteAsync(repository); // 根据success状态做后续处理,比如发送通知、记录日志 } }
3. DI处理多数据源实例的方案
ASP.NET Core原生DI默认只能注册一个接口的单个实例,而你有多个实现同一接口的仓库(比如两个IProductAttributeRepository),这里有两个常用的解决方案:
方案1:自定义仓库工厂(灵活通用)
创建一个工厂类,用来根据数据源名称获取对应的仓库实例:
public interface IRepositoryFactory { TRepo GetRepository<TRepo>(string dataSourceName) where TRepo : class; } public class RepositoryFactory : IRepositoryFactory { private readonly IServiceProvider _serviceProvider; public RepositoryFactory(IServiceProvider serviceProvider) { _serviceProvider = serviceProvider; } public TRepo GetRepository<TRepo>(string dataSourceName) where TRepo : class { return dataSourceName switch { "DataSource1" => _serviceProvider.GetRequiredService<DataSource1ProductAttributeRepository>() as TRepo, "DataSource2" => _serviceProvider.GetRequiredService<DataSource2ProductAttributeRepository>() as TRepo, "DataSource3" => _serviceProvider.GetRequiredService<DataSource3AssetRepository>() as TRepo, _ => throw new ArgumentException($"Unknown data source: {dataSourceName}") }; } }
然后在Program.cs中注册:
// 注册具体仓库类 builder.Services.AddScoped<DataSource1ProductAttributeRepository>(); builder.Services.AddScoped<DataSource2ProductAttributeRepository>(); builder.Services.AddScoped<DataSource3AssetRepository>(); // 注册工厂 builder.Services.AddScoped<IRepositoryFactory, RepositoryFactory>(); // 注册管道和步骤 builder.Services.AddScoped<IPipelineStep<IRepository<ProductAttributeEntity>, IEnumerable<ProductAttributeEntity>>, DataFetchStep<ProductAttributeEntity>>(); builder.Services.AddScoped<IPipelineStep<IEnumerable<ProductAttributeEntity>, IEnumerable<ProductAttributeEntity>>, DataValidationStep<ProductAttributeEntity>>(); builder.Services.AddScoped<IPipelineStep<IEnumerable<ProductAttributeEntity>, IEnumerable<UnifiedProductEntity>>, ProductToUnifiedEntityStep>(); builder.Services.AddScoped<IPipelineStep<IEnumerable<UnifiedProductEntity>, bool>, UnifiedDataPersistStep>(); // 组装管道 builder.Services.AddScoped<Pipeline<IRepository<ProductAttributeEntity>, bool>>(sp => new Pipeline<IRepository<ProductAttributeEntity>, bool>(new List<IPipelineStep<object, object>> { sp.GetRequiredService<IPipelineStep<IRepository<ProductAttributeEntity>, IEnumerable<ProductAttributeEntity>>>(), sp.GetRequiredService<IPipelineStep<IEnumerable<ProductAttributeEntity>, IEnumerable<ProductAttributeEntity>>>(), sp.GetRequiredService<IPipelineStep<IEnumerable<ProductAttributeEntity>, IEnumerable<UnifiedProductEntity>>>(), sp.GetRequiredService<IPipelineStep<IEnumerable<UnifiedProductEntity>, bool>>() }) );
使用时,通过工厂获取对应数据源的仓库:
public class DataCollectorService { private readonly IRepositoryFactory _repoFactory; private readonly ProductDataProcessingService _productProcessingService; public DataCollectorService(IRepositoryFactory repoFactory, ProductDataProcessingService productProcessingService) { _repoFactory = repoFactory; _productProcessingService = productProcessingService; } public async Task CollectAllDataAsync() { // 处理DataSource1的产品属性 var ds1Repo = _repoFactory.GetRepository<IProductAttributeRepository>("DataSource1"); await _productProcessingService.ProcessProductDataAsync(ds1Repo); // 处理DataSource2的产品属性 var ds2Repo = _repoFactory.GetRepository<IProductAttributeRepository>("DataSource2"); await _productProcessingService.ProcessProductDataAsync(ds2Repo); // 处理DataSource3的资产数据(同理可以创建资产数据的管道) var ds3Repo = _repoFactory.GetRepository<IAssetRepository>("DataSource3"); // await assetProcessingService.ProcessAssetDataAsync(ds3Repo); } }
方案2:标记接口(简单直接)
如果你的数据源数量不多且固定,可以给每个仓库类加一个标记接口,然后通过DI直接注入标记接口的实例:
// 定义标记接口 public interface IDataSource1Repo { } public interface IDataSource2Repo { } // 修改仓库类实现 public class DataSource1ProductAttributeRepository : IProductAttributeRepository, IDataSource1Repo { // 原有代码不变 } public class DataSource2ProductAttributeRepository : IProductAttributeRepository, IDataSource2Repo { // 原有代码不变 }
注册时:
builder.Services.AddScoped<IProductAttributeRepository, DataSource1ProductAttributeRepository>(); builder.Services.AddScoped<IProductAttributeRepository, DataSource2ProductAttributeRepository>(); // 同时注册标记接口 builder.Services.AddScoped<IDataSource1Repo, DataSource1ProductAttributeRepository>(); builder.Services.AddScoped<IDataSource2Repo, DataSource2ProductAttributeRepository>();
使用时直接注入标记接口:
public class DataCollectorService { private readonly IDataSource1Repo _ds1Repo; private readonly IDataSource2Repo _ds2Repo; private readonly ProductDataProcessingService _productProcessingService; public DataCollectorService(IDataSource1Repo ds1Repo, IDataSource2Repo ds2Repo, ProductDataProcessingService productProcessingService) { _ds1Repo = ds1Repo; _ds2Repo = ds2Repo; _productProcessingService = productProcessingService; } public async Task CollectProductAttributesAsync() { await _productProcessingService.ProcessProductDataAsync((IProductAttributeRepository)_ds1Repo); await _productProcessingService.ProcessProductDataAsync((IProductAttributeRepository)_ds2Repo); } }
总结
- 仓库模式帮你隔离了不同数据源的细节,让业务代码更干净,后续新增数据源只需要加新的仓库实现,不用修改上层逻辑。
- 管道模式把数据处理流程拆成独立步骤,你可以灵活添加、移除或调整步骤顺序,比如以后要加数据脱敏、缓存步骤,直接加个管道处理器就行。
- 多数据源的DI问题,工厂模式适合数据源较多、需要动态切换的场景;标记接口适合数据源固定、数量少的场景,选哪个看你的实际需求。
内容的提问来源于stack exchange,提问作者M. Haider
相关产品推荐
相关产品推荐

