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

如何用C#的Automatonymous结合RabbitMQ实现多应用状态机

Automatonymous + RabbitMQ 状态机实现指导

看起来你已经搭好了状态机的核心框架,我帮你把剩下的部分补全,并且梳理清楚四个独立应用的具体实现步骤——这样你就能快速复刻这个简单的请求解析流程了。

1. 先补全缺失的核心类型

你没贴出InterpreterInstance和相关的消息接口,这是状态机运行的基础,先把这些补上:

// 状态机实例,保存请求的状态和数据
public class InterpreterInstance : SagaStateMachineInstance
{
    public Guid CorrelationId { get; set; }
    public string CurrentState { get; set; }
    public Request Request { get; set; }
}

// 请求数据模型
public class Request
{
    public string RequestString { get; set; }
    public bool IsValid { get; set; }
    public string Answer { get; set; }

    public Request(string requestString)
    {
        RequestString = requestString;
    }
}

// 消息接口定义
public interface IRequesting
{
    Request Request { get; }
}

public interface IValidating
{
    Request Request { get; }
}

public interface IInterpreting
{
    Request Request { get; }
    string Answer { get; }
}

public interface IValidationNeeded
{
    Guid RequestId { get; }
    Request Request { get; }
}

public interface IInterpretationNeeded
{
    Guid RequestId { get; }
}

public interface IAnswerReady
{
    Guid RequestId { get; }
}

2. 补全状态机服务(RequestService)的完整代码

你的RequestService代码被截断了,这里是完整的实现,用来托管状态机并连接RabbitMQ:

using MassTransit;
using MassTransit.AutomatonymousIntegration;
using Topshelf;

public class RequestService : ServiceControl
{
    readonly IScheduler _scheduler;
    IBusControl _busControl;
    BusHandle _busHandle;
    InterpreterStateMachine _machine;
    InMemorySagaRepository<InterpreterInstance> _repository;

    public RequestService()
    {
        _scheduler = CreateScheduler();
    }

    private IScheduler CreateScheduler()
    {
        return new DefaultScheduler(new SystemClock());
    }

    public bool Start(HostControl hostControl)
    {
        Console.WriteLine("Creating bus...");
        
        _machine = new InterpreterStateMachine();
        _repository = new InMemorySagaRepository<InterpreterInstance>();

        // 配置RabbitMQ总线
        _busControl = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost/"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });

            // 注册状态机到接收端点
            cfg.ReceiveEndpoint(host, "interpreter_state_machine", e =>
            {
                e.StateMachineSaga(_machine, _repository);
            });
        });

        _busHandle = _busControl.Start();
        Console.WriteLine("State machine service started.");
        return true;
    }

    public bool Stop(HostControl hostControl)
    {
        _busHandle?.Stop();
        _busControl?.Stop();
        Console.WriteLine("State machine service stopped.");
        return true;
    }
}

// 服务启动入口
public class Program
{
    public static void Main()
    {
        HostFactory.Run(x =>
        {
            x.Service<RequestService>(s =>
            {
                s.ConstructUsing(() => new RequestService());
                s.WhenStarted(service => service.Start(null));
                s.WhenStopped(service => service.Stop(null));
            });

            x.RunAsLocalSystem();
            x.SetServiceName("InterpreterStateMachineService");
            x.SetDisplayName("Interpreter State Machine Service");
            x.SetDescription("Automatonymous State Machine for Request Interpretation");
        });
    }
}

3. 实现请求发送器应用

这个应用负责发送初始请求(比如"x=5")到状态机:

using MassTransit;

public class RequestSender
{
    static async Task Main(string[] args)
    {
        var bus = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            cfg.Host(new Uri("rabbitmq://localhost/"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });
        });

        await bus.StartAsync();
        try
        {
            // 发送请求消息
            await bus.Publish<IRequesting>(new
            {
                Request = new Request("x=5")
            });

            Console.WriteLine("Request 'x=5' sent successfully.");
            Console.ReadLine();
        }
        finally
        {
            await bus.StopAsync();
        }
    }
}

4. 实现请求验证器应用

这个应用监听IValidationNeeded事件,验证请求是否包含"=",然后发送IValidating消息回状态机:

using MassTransit;

public class RequestValidator : IConsumer<IValidationNeeded>
{
    public async Task Consume(ConsumeContext<IValidationNeeded> context)
    {
        var request = context.Message.Request;
        Console.WriteLine($"Validating request: {request.RequestString}");

        // 核心验证逻辑:检查请求字符串是否包含"="
        request.IsValid = request.RequestString.Contains("=");

        // 将验证结果发送回状态机
        await context.Publish<IValidating>(new
        {
            Request = request
        });

        Console.WriteLine($"Validation result for '{request.RequestString}': {request.IsValid}");
    }
}

public class ValidatorProgram
{
    static async Task Main(string[] args)
    {
        var bus = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost/"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });

            // 注册消费者到专属队列
            cfg.ReceiveEndpoint(host, "request_validator", e =>
            {
                e.Consumer<RequestValidator>();
            });
        });

        await bus.StartAsync();
        Console.WriteLine("Request validator service started.");
        Console.ReadLine();
        await bus.StopAsync();
    }
}

5. 实现请求解释器应用

这个应用监听IInterpretationNeeded事件,解析请求中的值,然后发送IInterpreting消息回状态机:

using MassTransit;

public class RequestInterpreter : IConsumer<IInterpretationNeeded>
{
    public async Task Consume(ConsumeContext<IInterpretationNeeded> context)
    {
        // 生产环境建议通过CorrelationId从状态机存储中查询请求数据
        // 这里为简化演示,直接模拟请求数据
        var request = new Request("x=5") { IsValid = true };
        var answer = request.RequestString.Split('=')[1];

        Console.WriteLine($"Interpreting request: {request.RequestString}, result: {answer}");

        // 将解析结果发送回状态机
        await context.Publish<IInterpreting>(new
        {
            Request = request,
            Answer = answer
        });
    }
}

public class InterpreterProgram
{
    static async Task Main(string[] args)
    {
        var bus = Bus.Factory.CreateUsingRabbitMq(cfg =>
        {
            var host = cfg.Host(new Uri("rabbitmq://localhost/"), h =>
            {
                h.Username("guest");
                h.Password("guest");
            });

            // 注册消费者到专属队列
            cfg.ReceiveEndpoint(host, "request_interpreter", e =>
            {
                e.Consumer<RequestInterpreter>();
            });
        });

        await bus.StartAsync();
        Console.WriteLine("Request interpreter service started.");
        Console.ReadLine();
        await bus.StopAsync();
    }
}

6. 关键注意事项

  • CorrelationId优化:你当前用RequestString作为关联键,测试场景没问题,但生产环境建议用唯一Guid作为CorrelationId,避免重复请求的状态冲突。
  • 状态持久化:示例用了InMemorySagaRepository,服务重启后状态会丢失,生产环境建议换成数据库存储(比如EntityFrameworkSagaRepository)。
  • 可靠性配置:如果需要保证消息不丢失,可以开启RabbitMQ的队列持久化、消息确认机制,以及MassTransit的重试策略。
  • 错误处理扩展:你定义了Error状态,可以补充错误事件的处理逻辑,比如发送告警通知,或者触发重试机制。

现在你可以按顺序启动四个应用:先启动状态机服务,再启动验证器、解释器,最后运行请求发送器,就能看到整个流程的完整执行了。

内容的提问来源于stack exchange,提问作者quorti

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:40:02