【问题标题】:How to create a generic pipeline in C#? [closed]如何在 C# 中创建通用管道? [关闭]
【发布时间】:2018-11-12 19:50:00
【问题描述】:

我正在尝试在 C# 中创建一个通用(通用)pipeline,以便在许多项目中重复使用。这个想法与ASP.NET Core Middleware 非常相似。它更像是一个可以动态组合的巨大函数(双向管道)(类似于BRE)

它需要获取一个输入模型,管理一系列之前加载的处理器,并在输入旁边返回一个封装在超模型中的输出模型。

这就是我所做的。我创建了一个Context 类,代表整体数据/模型:

public class Context<InputType, OutputType> where InputType : class, new() where OutputType : class, new()
{
    public Context()
    {
        UniqueToken = new Guid();
        Logs = new List<string>();
    }

    public InputType Input { get; set; } 

    public OutputType Output { get; set; }

    public Guid UniqueToken { get; }

    public DateTime ProcessStartedAt { get; set; }

    public DateTime ProcessEndedAt { get; set; }

    public long ProcessTimeInMilliseconds
    {
        get
        {
            return (long)ProcessEndedAt.Subtract(ProcessStartedAt).TotalMilliseconds;
        }
    }

    public List<string> Logs { get; set; }
}

然后我创建了一个接口,在真实处理器上强制签名:

public interface IProcessor
{
    void Process<InputType, OutputType>(Context<InputType, OutputType> context, IProcessor next) where InputType : class, new() where OutputType : class, new();
}

然后我创建了一个Container,来管理整个管道:

public class Container<InputType, OutputType> where InputType : class, new() where OutputType : class, new()
{
    public static List<IProcessor> Processors { get; set; }

    public static void Initialize()
    {
        LoadProcessors();
    }

    private static void LoadProcessors()
    {
        // loading processors from assemblies dynamically
    }

    public static Context<InputType, OutputType> Execute(InputType input)
    {
        if (Processors.Count == 0)
        {
            throw new FrameworkException("No processor is found to be executed");
        }
        if (input.IsNull())
        {
            throw new BusinessException($"{nameof(InputType)} is not provided for processing pipeline");
        }
        var message = new Context<InputType, OutputType>();
        message.Input = input;
        message.ProcessStartedAt = DateTime.Now;
        Processors[0].Process(message, Processors[1]);
        message.ProcessEndedAt = DateTime.Now;
        return message;
    }
}

我知道如何从给定文件夹中的程序集中动态加载处理器,所以这不是问题。但我被困在这些点上:

  1. 如何将下一个处理器注入每个处理器(我可以在每个处理器上强制使用Next 属性,但我想这是针对SRP,因为每个处理器应该只关心完成其工作,而不是保持链)
  2. 如何确保正确排序(一种选择是在每个处理器中都有Order属性,并确保它们没有重复值,但这似乎违反了SRP,每个处理器应该只关心关于处理,而不是关于它的顺序)
  3. 如何保证简单使用?在为团队创建基础架构时,对开发人员友好是一件大事。否则团队成员不会接受。
  4. 如何设计链条以使短路成为可能?

【问题讨论】:

  • 请编辑问题以将其限制为具有足够详细信息的特定问题,以确定适当的答案。避免一次问多个不同的问题。请参阅How to Ask 页面以获得澄清此问题的帮助。
  • 这是一个有趣的问题,我必须以非常相似的方式为插件和链接创建一个架构。但是,我不确定任何人都能理解您的意思,或者您真正想要实现的目标......没有模式,没有库,这只是设计架构,您是架构师.此外,任何建议都只是一种意见。我实际上认为您有能力自己解决这个问题并最了解您的担忧
  • 我认为容器负责流动。处理器要么在列表中排序要么每个处理器实例都应该用包含流的信息(如下一个属性)包装并存储在将由容器执行的列表中
  • 也许对 tpl 数据流感兴趣? (example)

标签: c# pipeline


【解决方案1】:

我会建议稍微不同的设计。这个想法基于装饰器模式。

首先,我将Context 设为非泛型类并删除输入和输出值。在我的设计中,上下文只包含上下文信息(如处理时间和消息):

public class Context
{
    public Context()
    {
        UniqueToken = new Guid();
        Logs = new List<string>();
    }        

    public Guid UniqueToken { get; }

    public DateTime ProcessStartedAt { get; set; }

    public DateTime ProcessEndedAt { get; set; }

    public long ProcessTimeInMilliseconds
    {
        get
        {
            return (long)ProcessEndedAt.Subtract(ProcessStartedAt).TotalMilliseconds;
        }
    }

    public List<string> Logs { get; set; }
}

然后,我会让处理器接口通用:

public interface IProcessor<InputType, OutputType>
{
    OutputType Process(InputType input, Context context);
}

然后我将您的 Container 变成了带有泛型类型参数的 Pipeline:

public interface IPipeline<InputType, OutputType>
{
    OutputType Execute(InputType input, out Context context);
    OutputType ExecuteSubPipeline(InputType input, Context context);
}

这两个函数的区别在于前者初始化上下文,后者只使用它。如果您不希望您的客户访问ExecuteSubPipeline(),您可能希望将其拆分为公共接口和内部接口。

然后的想法是将多个管道对象封装在彼此内部,这些对象具有越来越多的处理器。你从一个只有一个处理器的管道对象开始。比你将它包装在另一个管道对象中等等。为此,我从一个抽象基类开始。这个基类与一个处理器相关联,并且有一个函数AppendProcessor(),它创建一个添加了给定处理器的新管道:

public abstract class PipelineBase<InputType, ProcessorInputType, OutputType> : IPipeline<InputType, OutputType>
{
    protected IProcessor<ProcessorInputType, OutputType> currentProcessor;

    public PipelineBase(IProcessor<ProcessorInputType, OutputType> processor)
    {
        currentProcessor = processor;
    }

    public IPipeline<InputType, ProcessorOutputType> AppendProcessor<ProcessorOutputType>(IProcessor<OutputType, ProcessorOutputType> processor)
    {
        return new Pipeline<InputType, OutputType, ProcessorOutputType>(processor, this);
    }

    public OutputType Execute(InputType input, out Context context)
    {
        context = new Context();
        context.ProcessStartedAt = DateTime.Now;
        var result = ExecuteSubPipeline(input, context);
        context.ProcessEndedAt = DateTime.Now;
        return result;
    }

    public abstract OutputType ExecuteSubPipeline(InputType input, Context context);
}

现在,我们有这个管道的两个具体实现:一个是任何管道的起点的终端实现和一个包装器管道:

public class TerminalPipeline<InputType, OutputType> : PipelineBase<InputType, InputType, OutputType>
{       
    public TerminalPipeline(IProcessor<InputType, OutputType> processor)
        :base(processor)
    { }

    public override OutputType ExecuteSubPipeline(InputType input, Context context)
    {
        return currentProcessor.Process(input, context);
    }
}

public class Pipeline<InputType, ProcessorInputType, OutputType> : PipelineBase<InputType, ProcessorInputType, OutputType>
{
    IPipeline<InputType, ProcessorInputType> previousPipeline;

    public Pipeline(IProcessor<ProcessorInputType, OutputType> processor, IPipeline<InputType, ProcessorInputType> previousPipeline)
        : base(processor)
    {
        this.previousPipeline = previousPipeline;
    }

    public override OutputType ExecuteSubPipeline(InputType input, Context context)
    {
        var previousPipelineResult = previousPipeline.ExecuteSubPipeline(input, context);
        return currentProcessor.Process(previousPipelineResult, context);
    }
}

为了方便使用,我们还创建一个辅助函数,用于创建终端启动管道(以允许类型参数推导):

public static class Pipeline
{
    public static TerminalPipeline<InputType, OutputType> Create<InputType, OutputType>(IProcessor<InputType, OutputType> processor)
    {
        return new TerminalPipeline<InputType, OutputType>(processor);
    }
}

然后,我们可以将这种结构与各种处理器一起使用。例如:

class FloatToStringProcessor : IProcessor<float, string>
{
    public string Process(float input, Context context)
    {
        return input.ToString();
    }
}

class RepeatStringProcessor : IProcessor<string, string>
{
    public string Process(string input, Context context)
    {
        return input + input + input;
    }
}

class Program
{
    public static void Main()
    {
        var pipeline = Pipeline
            .Create(new FloatToStringProcessor())
            .AppendProcessor(new RepeatStringProcessor());

        Context ctx;
        var result = pipeline.Execute(5, out ctx);
        Console.WriteLine($"Pipeline result: {result}");
        Console.WriteLine($"Pipeline execution took {ctx.ProcessTimeInMilliseconds} milliseconds");
    }
}

这将打印出来

Pipeline result: 555
Pipeline execution took 6 milliseconds

我不明白你所说的短路是什么意思。在我看来,短路仅对(至少)不需要评估一个操作数的二元运算符有意义。但是由于您的运算符都是一元的,因此不能真正应用。处理器可以随时检查输入,当发现不需要处理时直接返回。

可以通过在IPipeline 接口中添加类似LoadProcessors() 的内容轻松添加动态加载,类似于ExecuteSubPipeline()。在这种情况下,处理器对象必须是代表(仍然正确键入)。然后,LoadProcessors() 可以在加载它们后将它们替换为实际的处理器。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-12-29
    • 1970-01-01
    • 1970-01-01
    • 2019-03-15
    • 2011-09-18
    • 2016-03-28
    • 2015-08-11
    相关资源
    最近更新 更多