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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 07:22:37