【问题标题】:How to do WaitAll with Akka.Net?如何用 Akka.Net 做 WaitAll?
【发布时间】:2016-02-16 14:12:59
【问题描述】:

我在 Akka.Net 中有一个演员层次结构,我想知道我是否选择了正确的方法来做某事,或者是否有更好/更简单的方法来实现我想要的。

我的具体示例是,我正在构造一个 User 演员以响应用户登录系统,并且在构造该演员时,我需要两条数据来完成演员的建设。

如果这是常规的 .NET 代码,我可能会有类似以下内容...

public Task<User> LoadUserAsync (string username)
{
  IProfileService profileService = ...;
  IMessageService messageService = ...;

  var loadProfileTask = profileService.GetUserProfileAsync(username);
  var loadMessagesTask = messageService.GetMessagesAsync(username);

  Task.WaitAll(loadProfileTask, loadMessagesTask);

  // Now construct the user from the result of both tasks
  var user = new User
  {
    Profile = loadProfileTask.Result,
    Messages = loadMessagesTask.Result
  }

  return Task.FromResult(user);
}

这里我使用WaitAll等待下级任务完成,让它们并发运行。

我的问题是 - 如果我想在 Akka.Net 中做同样的事情,以下是最常见的方法吗?如图所示,我创建了以下...

当我创建我的用户角色时,我会构建一个(临时)用户加载器角色,它的工作是通过调用配置文件角色和消息角色来获取完整的用户详细信息。获取数据的叶子actor如下...

public class UserProfileLoader : ReceiveActor
{
    public UserProfileLoader()
    {
        Receive<LoadUserRequest>(msg =>
        {
            // Load the user profile from somewhere
            var profile = new UserProfile();

            // And respond to the Sender
            Sender.Tell(profile);
            Self.Tell(PoisonPill.Instance);
        });
    }
}

public class UserMessagesLoader : ReceiveActor
{
    public UserMessagesLoader()
    {
        Receive<LoadUserRequest>(msg =>
        {
            // Load the messages from somewhere
            var messages = new List<Message>();

            // And respond to the Sender
            Sender.Tell(messages);
            Self.Tell(PoisonPill.Instance);
        });
    }
}

他们从哪里得到这个讨论的数据并不重要,但两者都只是通过返回一些数据来响应请求。

然后我有协调两个数据收集参与者的参与者......

public class UserLoaderActor : ReceiveActor
{
    public UserLoaderActor()
    {
        Receive<LoadUserRequest>(msg => LoadProfileAndMessages(msg));
        Receive<UserProfile>(msg =>
        {
            _profile = msg;
            FinishIfPossible();
        });

        Receive<List<Message>>(msg =>
        {
            _messages = msg;
            FinishIfPossible();
        });
    }

    private void LoadProfileAndMessages(LoadUserRequest msg)
    {
        _originalSender = Sender;
        Context.ActorOf<UserProfileLoader>().Tell(msg);
        Context.ActorOf<UserMessagesLoader>().Tell(msg);
    }

    private void FinishIfPossible()
    {
        if ((null != _messages) && (null != _profile))
        {
            _originalSender.Tell(new LoadUserResponse(_profile, _messages));
            Self.Tell(PoisonPill.Instance);
        }
    }

    private IActorRef _originalSender;
    private UserProfile _profile;
    private List<Message> _messages;
}

这只是创建了两个从属参与者,向他们发送消息以进行破解,然后等待两者都响应,然后再将收集到的所有数据发送回原始请求者。

那么,这似乎是一种合理的方式来协调两个不同的响应,以便将它们组合起来?有没有比自己制作更简单的方法?

提前感谢您的回复!

【问题讨论】:

  • 我不知道您是否需要单独的演员作为只收到 2 条回复的门面。 Useractor 不能直接发送这两个请求吗?
  • 我确实可以让用户actor完成所有这些工作,但是由于用户actor很可能有更多的事情要做,所以我选择将这里的责任分配给加载器.不管工作分配如何,上述内容是否有意义和/或是否有更好/更少代码/更标准的方法来做到这一点?
  • 配置文件和消息参与者中是否有任何非常具体的事情发生,或者他们只是从数据库中获取数据?如果是这样,我可能会将其包装在一些异步方法中,该方法在两个子任务上执行WhenAll,然后返回该任务,以便调用者在完成后可以PipeTo(Self)。我知道我们宣扬“将危险的工作推给孩子”但是需要考虑您是否真的从中受益,或者它是否只会导致代码膨胀..

标签: akka.net


【解决方案1】:

谢谢大家,所以我现在根据 Roger 和 Jeff 的建议将演员简化为以下内容...

public class TaskBasedUserLoader : ReceiveActor
{
    public TaskBasedUserLoader()
    {
        Receive<LoadUserRequest>(msg => LoadProfileAndMessages(msg));
    }

    private void LoadProfileAndMessages(LoadUserRequest msg)
    {
        var originalSender = Sender;
        var loadPreferences = this.LoadProfile(msg.UserId);
        var loadMessages = this.LoadMessages(msg.UserId);

        Task.WhenAll(loadPreferences, loadMessages)
            .ContinueWith(t => new UserLoadedResponse(loadPreferences.Result, loadMessages.Result), 
                TaskContinuationOptions.AttachedToParent & TaskContinuationOptions.ExecuteSynchronously)
            .PipeTo(originalSender);
    }

    private Task<UserProfile> LoadProfile(string userId)
    {
        return Task.FromResult(new UserProfile { UserId = userId });
    }

    private Task<List<Message>> LoadMessages(string userId)
    {
        return Task.FromResult(new List<Message>());
    }
}

LoadProfile 和 LoadMessages 方法最终将调用存储库来获取数据,但现在我有一个简洁的方法来做我想做的事。

再次感谢!

【讨论】:

  • 很高兴为您提供帮助。注意 AttachedToParent,这里的“父”不是 Task.WhenAll,它是调用 LoadProfileAndMessages 的当前任务(在这种情况下没有)。而你使用 & 而不是 |组合标志,所以你最终得到零。
  • 顺便说一下,我对 TaskContinuationOptions 的使用是从 Petabridge 关于在此处使用 PipeTo 的文章中复制的 - petabridge.com/blog/akkadotnet-async-actors-using-pipeto。所以我想这也需要更新,因为同样的问题也有两个实例。
【解决方案2】:

恕我直言,这是一个有效的过程,因为您分叉操作然后加入它。

顺便说一句,您可以使用this.Self.GracefulStop(new TimeSpan(1)); 而不是发送毒丸。

【讨论】:

【解决方案3】:

您可以结合使用 Ask、WhenAll 和 PipeTo:

var task1 = actor1.Ask<Result1>(request1);
var task2 = actor2.Ask<Result2>(request2);

Task.WhenAll(task1, task2)
    .ContinueWith(_ => new Result3(task1.Result, task2.Result))
    .PipeTo(Self);

...

Receive<Result3>(msg => { ... });

【讨论】:

    猜你喜欢
    • 2011-09-01
    • 1970-01-01
    • 2015-07-13
    • 1970-01-01
    • 2014-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多