【问题标题】:How to get state from service fabric actor without waiting for other methods to complete?如何在不等待其他方法完成的情况下从服务结构参与者获取状态?
【发布时间】:2019-04-24 05:49:00
【问题描述】:

我有一个正在运行的服务正在迭代 X 演员,使用 ActorProxy 询问他们的状态。

对我来说重要的是,此服务不会因等待来自 ect 提醒回调的 actor 中的其他一些长时间运行的方法而被阻塞。

是否有某种方法可以调用下面的简单示例 GetState(),这将允许该方法以正确的方式完成,而不会阻止某些提醒运行。

class Actor : IMyActor{

public Task<MyState> GetState() => StateManager.GetAsync("key")
}

另一种选择。

调用服务的正确方法是什么,如果它在 5 秒内没有回复,请包含。

var proxy = ActorProxy.Create<IMyActor();
var state = await proxy.GetState(); // this will wait until the actor is ready to send back the state. 

【问题讨论】:

  • 不要对actor执行长时间运行的动作——它应该能够立即响应。 -- 如果你需要一个长时间运行的动作,把它放到一个子actor上,当它完成时通知父actor。 -- 这样父actor可以保持响应。

标签: c# azure-service-fabric


【解决方案1】:

即使是当前正在执行阻塞方法的 Actor,也可以读取 Actor 状态。 Actors 使用IActorStateManager 存储他们的状态,而IActorStateProvider 又使用IActorStateProviderIActorStateProvider 每个 ActorService 实例化一次。每个分区都会实例化负责托管和运行参与者的ActorService。 Actor 服务的核心是StatefulService(或者更确切地说是StatefulServiceBase,它是常规有状态服务使用的基类)。考虑到这一点,我们可以像使用常规服务一样使用迎合我们 Actor 的 ActorService,即使用基于 IService 的服务接口。

IActorStateProvider(如果您使用的是持久状态,则由KvsActorStateProvider 实现)有两种我们可以使用的方法:

Task<T> LoadStateAsync<T>(ActorId actorId, string stateName, CancellationToken cancellationToken = null);
Task<PagedResult<ActorId>> GetActorsAsync(int numItemsToReturn, ContinuationToken continuationToken, CancellationToken cancellationToken);

对这些方法的调用不受参与者锁的影响,这是有道理的,因为它们旨在支持分区上的所有参与者。

示例:

创建一个自定义 ActorService 并使用它来托管您的演员:

public interface IManyfoldActorService : IService
{
    Task<IDictionary<long, int>> GetCountsAsync(CancellationToken cancellationToken);
}

public class ManyfoldActorService : ActorService, IManyfoldActorService
{
    ...
}

Program.Main注册新的ActorService:

ActorRuntime.RegisterActorAsync<ManyfoldActor>(
    (context, actorType) => new ManyfoldActorService(context, actorType)).GetAwaiter().GetResult();

假设我们有一个具有以下方法的简单 Actor:

    Task IManyfoldActor.SetCountAsync(int count, CancellationToken cancellationToken)
    {
        Task.Delay(TimeSpan.FromSeconds(30), cancellationToken).GetAwaiter().GetResult();
        var task = this.StateManager.SetStateAsync("count", count, cancellationToken);
        ActorEventSource.Current.ActorMessage(this, $"Finished set {count} on {this.Id.GetLongId()}");
        return task;
    }

它等待 30 秒(模拟长时间运行、阻塞、方法调用),然后将状态值 "count" 设置为 int

在一个单独的服务中,我们现在可以调用 SetCountAsync 让 Actors 生成一些状态数据:

    protected override async Task RunAsync(CancellationToken cancellationToken)
    {
        var actorProxyFactory = new ActorProxyFactory();
        long iterations = 0;
        while (true)
        {
            cancellationToken.ThrowIfCancellationRequested();
            iterations += 1;
            var actorId = iterations % 10;
            var count = Environment.TickCount % 100;
            var manyfoldActor = actorProxyFactory.CreateActorProxy<IManyfoldActor>(new ActorId(actorId));
            manyfoldActor.SetCountAsync(count, cancellationToken).ConfigureAwait(false);
            ServiceEventSource.Current.ServiceMessage(this.Context, $"Set count {count} on {actorId} @ {iterations}");
            await Task.Delay(TimeSpan.FromSeconds(3), cancellationToken);
        }
    }

此方法只是循环不断地更改演员的值。 (注意总共 10 个 Actor 之间的相关性,延迟 3 秒和 Actor 延迟 30 秒。简单地设计这种方式是为了防止等待锁定的 Actor 调用的无限累积)。每个调用也作为即发即弃的方式执行,因此我们可以在下一个actor返回之前继续更新下一个actor的状态。这是一段愚蠢的代码,它只是为了证明理论而设计的。

现在在actor服务中我们可以像这样实现GetCountsAsync方法:

    public async Task<IDictionary<long, int>> GetCountsAsync(CancellationToken cancellationToken)
    {
        ContinuationToken continuationToken = null;
        var actors = new Dictionary<long, int>();

        do
        {
            var page = await this.StateProvider.GetActorsAsync(100, continuationToken, cancellationToken);

            foreach (var actor in page.Items)
            {
                var count = await this.StateProvider.LoadStateAsync<int>(actor, "count", cancellationToken);
                actors.Add(actor.GetLongId(), count);
            }

            continuationToken = page.ContinuationToken;
        }
        while (continuationToken != null);

        return actors;
    }

这使用底层的ActorStateProvider 来查询所有已知的Actor(针对该分区),然后直接读取每个“绕过”Actor 并且不被Actor 的方法执行阻塞的状态。

最后一部分,一些可以调用我们的 ActorService 并在所有分区中调用 GetCountsAsync 的方法:

    public IDictionary<long, int> Get()
    {
        var applicationName = FabricRuntime.GetActivationContext().ApplicationName;
        var actorServiceName = $"{typeof(IManyfoldActorService).Name.Substring(1)}";
        var actorServiceUri = new Uri($"{applicationName}/{actorServiceName}");

        var fabricClient = new FabricClient();
        var partitions = new List<long>();
        var servicePartitionList = fabricClient.QueryManager.GetPartitionListAsync(actorServiceUri).GetAwaiter().GetResult();
        foreach (var servicePartition in servicePartitionList)
        {
            var partitionInformation = servicePartition.PartitionInformation as Int64RangePartitionInformation;
            partitions.Add(partitionInformation.LowKey);
        }

        var serviceProxyFactory = new ServiceProxyFactory();

        var actors = new Dictionary<long, int>();
        foreach (var partition in partitions)
        {
            var actorService = serviceProxyFactory.CreateServiceProxy<IManyfoldActorService>(actorServiceUri, new ServicePartitionKey(partition));

            var counts = actorService.GetCountsAsync(CancellationToken.None).GetAwaiter().GetResult();
            foreach (var count in counts)
            {
                actors.Add(count.Key, count.Value);
            }
        }
        return actors;
    }

运行此代码现在将为我们提供 10 个参与者,它们每 33:d 秒更新一次状态,并且每个参与者每次忙 30 秒。当每个 Actor 方法返回时,Actor 服务就会看到更新的状态。

此示例中省略了一些内容,例如,当您在 Actor 服务中加载状态时,我们可能应该防止超时。

【讨论】:

  • 很棒的答案:)
  • yoape - 我使用了这种方法,现在我发现有些奇怪。当使用 actorservice 上的 stateprovider 来更新某个 actor 的某些状态时。然后一个演员已经激活了加载状态,它得到的是旧值而不是新的更新值。这也是你看到的吗?
  • @PoulK.Sørensen,我们仍在生产代码中使用这种方法。我们还没有看到那里的这种行为,似乎对我们有用。您运行的是哪个版本的 SDK?
  • 这个特定的 2.7.198 - 也将很快在另一个项目上使用 latest 进行验证。我通过让演员在停用时向演员服务报告来解决这个问题。然后,如果参与者服务看到新状态可用,它将在停用后再次激活参与者,然后参与者具有最新状态。所以我看到的问题是从演员服务更新的状态在已经活跃的演员中没有更新。 (也许它的状态在内存中,并且参与者服务正在操作磁盘上的状态)
【解决方案2】:

没有办法做到这一点。 Actor 是单线程的。如果他们正在等待在任何actor方法中完成的长时间运行的工作,那么任何其他方法(包括来自外部的方法)都必须等待。

【讨论】:

  • 我们在早期的 sdk 版本中没有一些 ReadOnly 属性吗?有没有办法让给定参与者的状态管理器然后查找状态而无需等待。我确信必须有我想要的解决方案?
  • Readonly 在早期版本中确实存在,但它暗示了在方法退出后是否需要复制状态并且对单线程锁定语义没有影响。
【解决方案3】:

感谢您的所有帮助。能够以您的示例为例,并通过一些调整使其工作。我遇到的唯一问题是将数据传回原始应用程序服务时遇到了未知类型。得到了

“ArrayOfKeyValueOflonglong 不是预期的。将任何静态未知的类型添加到已知类型列表中 - 例如,通过使用 KnownTypeAttribute 属性或将它们添加到传递给 DataContractSerializer 的已知类型列表中”

所以我将 GetCountsAsync 的返回类型更改为 List ,并在我的类中使用了 DataContract 和 DataMember 类属性,它运行良好。似乎从分区中的许多参与者检索状态数据的能力应该是参与者服务的核心部分,您不应该创建自定义参与者服务来获取 StateProvider 信息。再次感谢您!

【讨论】:

  • 您能否编辑您的答案以包含您的代码,而不仅仅是描述它?
猜你喜欢
  • 2017-12-08
  • 2013-01-29
  • 1970-01-01
  • 2019-04-10
  • 1970-01-01
  • 2018-03-06
  • 1970-01-01
  • 2017-06-08
  • 1970-01-01
相关资源
最近更新 更多