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

Akka.NET+ASP.NET模块化单体架构:多Actor交互后向控制器返回结果

使用Akka.NET + ASP.NET Core实现REST服务的响应返回问题

我基于Akka.NET官方示例实现了一个REST服务,创建了包含FooActor引用的AkkaService,以及将HTTP请求转为RunProcess消息发送给FooActor的MyController。

现有代码如下:

控制器代码

[Route("api/[controller]")]
[ApiController]
public class MyController : Controller
{
    private readonly ILogger<MyController> _logger;
    private readonly IAkkaService Service;

    public RebalancingController(ILogger<MyController> logger, IAkkaService bridge)
    {
        _logger = logger;
        Service = bridge;
    }

    [HttpGet]
    public async Task<ProcessTerminated> Get()
    {
        var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60));
        return await Service.RunProcess(cts.Token);
    }
}

AkkaService代码

public class AkkaService : IAkkaService, IHostedService
{
    private ActorSystem ActorSystem { get; set; }
    public IActorRef FooActor { get; private set; }
    private readonly IServiceProvider ServiceProvider;

    public AkkaService(IServiceProvider sp)
    {
        ServiceProvider = sp;
    }

    public async Task StartAsync(CancellationToken cancellationToken)
    {
        var hocon = ConfigurationFactory.ParseString(await File.ReadAllTextAsync("app.conf", cancellationToken));
        var bootstrap = BootstrapSetup.Create().WithConfig(hocon);
        var di = DependencyResolverSetup.Create(ServiceProvider);
        var actorSystemSetup = bootstrap.And(di);
        ActorSystem = ActorSystem.Create("AkkaSandbox", actorSystemSetup);
        
        // 通过依赖注入创建Props
        var fooProps = DependencyResolver.For(ActorSystem).Props<FooActor>();
        FooActor = ActorSystem.ActorOf(fooProps.WithRouter(FromConfig.Instance), "foo");

        await Task.CompletedTask;
    }

    public async Task<ProcessTerminated> RunProcess(CancellationToken token)
    {
        return await FooActor.Ask<ProcessTerminated>(new RunProcess(), token);
    }
}

FooActor代码

public FooActor(IServiceProvider sp)
{
    _scope = sp.CreateScope();

    Receive<RunProcess>(x =>
    {
        var barActor = Context.ActorOf(Props.Create<BarActor>(sp), "BarActor");
        barActor.Tell(new BarRequest());
        _log.Info("Sending a request to Bar Actor ");
    });

    Receive<BarResponse>(x =>
    {
        // 这里需要向控制器返回ProcessTerminated消息
    });
}

当前场景与需求

场景:FooActor向BarActor发送BarRequest消息并等待BarResponse,之后需要向控制器返回ProcessTerminated消息。

要求:

  • 解耦BarActor与FooActor:BarActor不需要知道FooActor或控制器的存在,消息结构和响应逻辑不能依赖FooActor的后续处理;
  • 支持复杂多Actor交互:最终结果能正确返回给发起请求的控制器。

解决方案

核心思路:利用Akka.NET的请求上下文关联机制

当FooActor处理RunProcess消息时,保留原始请求的发送者(即Ask的发起方),通过临时Actor跟踪单个请求的生命周期,在收到BarResponse后将ProcessTerminated回复给原始发送者,同时保证BarActor的完全解耦。

修改后的FooActor实现

public FooActor(IServiceProvider sp)
{
    _scope = sp.CreateScope();

    Receive<RunProcess>(x =>
    {
        // 保存原始请求的发送者(即AkkaService中Ask的调用方)
        var originalSender = Sender;
        var barActor = Context.ActorOf(Props.Create<BarActor>(sp), $"BarActor-{Guid.NewGuid()}");
        
        // 创建临时Actor处理当前请求的后续响应,隔离不同请求的状态
        var requestHandler = Context.ActorOf(Props.Create(() => new RequestHandlerActor(originalSender)));
        
        // 监听BarActor生命周期,处理异常终止情况
        Context.Watch(barActor);
        // 将BarRequest转发给BarActor,让BarActor直接回复给临时处理Actor
        barActor.Forward(new BarRequest());
    });
}

// 临时请求处理Actor:负责单个请求的响应转发与生命周期管理
public class RequestHandlerActor : ReceiveActor
{
    private readonly IActorRef _originalSender;

    public RequestHandlerActor(IActorRef originalSender)
    {
        _originalSender = originalSender;

        Receive<BarResponse>(x =>
        {
            // 收到BarResponse后,向原始发送者回复ProcessTerminated
            _originalSender.Tell(new ProcessTerminated());
            // 处理完成后停止自身
            Context.Stop(Self);
        });

        // 处理BarActor意外终止的异常场景
        Receive<Terminated>(x =>
        {
            // 可自定义错误消息返回给控制器
            _originalSender.Tell(new ProcessFailed("BarActor terminated unexpectedly"));
            Context.Stop(Self);
        });
    }
}

BarActor的解耦实现

public class BarActor : ReceiveActor
{
    public BarActor()
    {
        Receive<BarRequest>(x =>
        {
            // 执行业务逻辑,生成BarResponse
            var response = new BarResponse();
            // 直接回复给消息发送者(即临时RequestHandlerActor),无需知晓后续处理
            Sender.Tell(response);
            // 一次性Actor处理完成后停止自身
            Context.Stop(Self);
        });
    }
}

关键说明

  1. 完全解耦:BarActor仅需将结果回复给消息发送者,完全不依赖FooActor或控制器的存在,符合解耦要求;
  2. 复杂场景适配:临时RequestHandlerActor为每个请求独立维护状态,后续若需添加更多Actor交互(如调用BazActor),可直接在该Actor中扩展逻辑,保证结果能正确回溯到控制器;
  3. 并发安全:通过临时Actor隔离不同请求的状态,避免在FooActor中维护复杂的请求映射表,降低并发冲突风险;
  4. 异常兜底:通过Context.Watch监听BarActor生命周期,当Actor意外终止时可向控制器返回错误信息,避免请求无限等待。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 02:06:25