如何用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
相关产品推荐
相关产品推荐

