Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Приложения взаимодействуют с потоками через API, подобные хорошо известным Reactive Extensions (Rx) в .NET. Основное различие заключается в том, что Orleans расширения потоков являются асинхронными для повышения эффективности обработки в Orleansраспределенной и масштабируемой вычислительной структуре.
Асинхронный поток
Начните с использования поставщика потоков для получения дескриптора к потоку. Поставщик потоков можно рассматривать как фабрику потоков, которая позволяет реализующим настраивать поведение потоков и семантику:
public async Task SetupStream()
{
IStreamProvider streamProvider = this.GetStreamProvider("SimpleStreamProvider");
StreamId streamId = StreamId.Create("MyStreamNamespace", this.GetPrimaryKey());
IAsyncStream<string> stream = streamProvider.GetStream<string>(streamId);
}
IStreamProvider streamProvider = base.GetStreamProvider("SimpleStreamProvider");
IAsyncStream<T> stream = streamProvider.GetStream<T>(Guid, "MyStreamNamespace");
Вы можете получить доступ к поставщику потоков либо путем вызова метода Grain.GetStreamProvider внутри зерна, либо путем вызова метода GetStreamProvider на экземпляре клиента.
Orleans.Streams.IAsyncStream<T> — это логический, строго типизированный дескриптор виртуального потока, аналогичный ссылке на Orleans (Grain Reference).
GetStreamProvider Вызовы и GetStream являются исключительно локальными. Аргументы для GetStream включают GUID и дополнительную строку, называемую пространством имен потока (которая может быть null). Вместе GUID и строка пространства имен составляют идентификатор потока (аналогично аргументам для IGrainFactory.GetGrain). Эта комбинация обеспечивает дополнительную гибкость при определении идентификации потоков. Как и зерно 7, которое может существовать в пределах PlayerGrain типа, и другое зерно 7, которое может существовать в пределах ChatRoomGrain типа, Stream 123 может существовать в PlayerEventsStream пространстве имен, а другой поток 123 может существовать в ChatRoomMessagesStream пространстве имен.
Производство и потребление
IAsyncStream<T> реализует как IAsyncObserver<T> интерфейсы, так и IAsyncObservable<T> интерфейсы. Это позволяет вашему приложению использовать поток либо для создания новых событий с помощью IAsyncObserver<T>, либо для подписки и потребления событий с помощью IAsyncObservable<T>.
public interface IAsyncObserver<in T>
{
Task OnNextAsync(T item, StreamSequenceToken token = null);
Task OnCompletedAsync();
Task OnErrorAsync(Exception ex);
}
public interface IAsyncObservable<T>
{
Task<StreamSubscriptionHandle<T>> SubscribeAsync(IAsyncObserver<T> observer);
}
Чтобы создать события в потоке, приложение вызывает:
await stream.OnNextAsync<T>(event)
Чтобы подписаться на поток, приложение вызывает:
StreamSubscriptionHandle<T> subscriptionHandle = await stream.SubscribeAsync(IAsyncObserver)
Аргумент SubscribeAsync может быть либо объектом, реализующим IAsyncObserver<T> интерфейс, либо сочетанием лямбда-функций для обработки входящих событий. Дополнительные варианты SubscribeAsync доступны через AsyncObservableExtensions класс. SubscribeAsync возвращает StreamSubscriptionHandle<T>, непрозрачный дескриптор, используемый для отмены подписки из потока (по аналогии с асинхронной версией IDisposable).
await subscriptionHandle.UnsubscribeAsync()
Важно отметить, что подписка предназначена для определённого элемента, а не для активации. Когда код зерна подписывается на поток, эта подписка превышает жизнь этой активации и остается постоянной до тех пор, пока код зерна (потенциально в другой активации) явно отменяет подписку. Это ядро абстракции виртуального потока: не только все потоки всегда существуют логически, но подписка на поток также устойчива и живет за пределами конкретной физической активации, созданной ею.
Кратность
Поток Orleans может содержать несколько производителей и нескольких потребителей. Сообщение, опубликованное производителем, доставляется всем потребителям, подписанным на поток до публикации сообщения.
Кроме того, потребитель может подписаться на один поток несколько раз. Каждый раз, когда он подписывается, он возвращает уникальный StreamSubscriptionHandle<T>. Если зерно (или клиент) подписывается X раз на один и тот же поток, оно получает одно и то же событие X раз, по одному за каждую подписку. Потребитель также может отменить подписку из отдельной подписки. Вы можете найти все текущие подписки, позвонив по номеру:
IList<StreamSubscriptionHandle<T>> allMyHandles =
await IAsyncStream<T>.GetAllSubscriptionHandles();
Восстановление после сбоев
Если производитель потока прекращает работу (или его зерно деактивировано), не нужно ничего делать. В следующий раз, когда это зерно хочет создать больше событий, он может снова получить дескриптор потока и создать новые события как обычно.
Логика потребителя немного более сложная. Как уже упоминалось ранее, после того как потребительский зерно подписывается на поток, эта подписка действует до тех пор, пока зерно не решит явно аннулировать её. Если потребитель потока умирает (или если его зерно деактивируется), и в потоке создается новое событие, зерно потребителя автоматически активируется (точно так же, как любое обычное Orleans зерно автоматически активируется при отправке ему сообщения). Единственное, что сейчас нужно сделать в коде зерна, — это предоставить IAsyncObserver<T> для обработки данных. Потребитель должен повторно подключить логику обработки в рамках OnActivateAsync() метода. Для этого он может вызвать:
StreamSubscriptionHandle<int> newHandle =
await subscriptionHandle.ResumeAsync(IAsyncObserver);
Потребитель использует предыдущий идентификатор, полученный во время начальной подписки, для возобновления обработки. Обратите внимание, что ResumeAsync просто обновляет существующую подписку с новым экземпляром логики IAsyncObserver<T> и не изменяет тот факт, что этот потребитель уже подписан на этот поток.
Как потребитель может получить старый subscriptionHandle? Существует два варианта. Возможно, потребитель сохранил дескриптор, возвращенный из исходной SubscribeAsync операции, и теперь может его использовать. Кроме того, если у потребителя нет дескриптора, он может попросить IAsyncStream<T> все его активные дескриптор подписки, вызвав:
IList<StreamSubscriptionHandle<T>> allMyHandles =
await IAsyncStream<T>.GetAllSubscriptionHandles();
Затем потребитель может возобновить все из них или отменить подписку от некоторых при необходимости.
Совет
Если потребитель интерфейса IAsyncObserver<T> реализует его напрямую (public class MyGrain<T> : Grain, IAsyncObserver<T>), то теоретически ему не потребуется повторно подключать IAsyncObserver<T>, поэтому вызов ResumeAsync также не потребуется. Среда выполнения потоковой передачи должна автоматически определить, что зерно уже содержит реализацию IAsyncObserver<T> и вызвать эти IAsyncObserver<T> методы. Однако среда выполнения потоковой передачи в настоящее время не поддерживает это, и код зерна по-прежнему должен явно вызывать ResumeAsync, даже если зерно реализует IAsyncObserver<T> напрямую.
Явные и неявные подписки
По умолчанию, для доступа к потоку, потребитель должен явно на него подписаться. Эта подписка обычно активируется внешним сообщением, которое зерно (или клиент) получает с указанием подписаться. Например, в службе чата, когда пользователь присоединяется к комнате чата, их grain получает сообщение с именем чата, из-за чего пользователь grain подписывается на этот поток чата.
Кроме того, Orleans потоки поддерживают неявные подписки. В этой модели зерно явно не подписывается. Это подписывается автоматически и неявно на основе его идентификатора зерна и ImplicitStreamSubscriptionAttribute. Основное преимущество неявных подписок заключается в том, что активность потока может автоматически активировать зерно (а следовательно, и подписку). Например, при использовании SMS-потоков, если одно зерно хотело производить поток, а другое зерно его обрабатывало, производителю понадобились бы идентификационные данные потребителя зерна, чтобы сделать вызов метода зерна, призывая его подписаться на поток. Только тогда он мог приступить к отправке событий. Вместо этого, при использовании неявных подписок, производитель может просто начать передавать события в поток, и потребительский модуль автоматически активируется и подключается. В этом случае продюсеру не нужно знать, кто читает события.
Реализация зерна MyGrainType может объявлять атрибут [ImplicitStreamSubscription("MyStreamNamespace")]. Это сообщает среде выполнения потоковой передачи, что, когда событие создается в потоке с идентификатором GUID XXX и пространством имен "MyStreamNamespace", оно должно быть доставлено в зерно с идентификатором XXX типа MyGrainType. То есть среда выполнения сопоставляет поток <XXX, MyStreamNamespace> с потребительским зерном <XXX, MyGrainType>.
Наличие ImplicitStreamSubscription среды выполнения потоковой передачи автоматически подписывает это зерно на поток и передает в нее события потоковой передачи. Однако код зерна по-прежнему должен сообщить среде выполнения, как она хочет обрабатывать события. По сути, оно должно быть присоединено IAsyncObserver<T>. Поэтому при активации зерна код зерна внутри OnActivateAsync должен вызываться:
public override async Task OnActivateAsync(CancellationToken cancellationToken)
{
IStreamProvider streamProvider =
this.GetStreamProvider("SimpleStreamProvider");
StreamId streamId =
StreamId.Create("MyStreamNamespace", this.GetPrimaryKey());
IAsyncStream<string> stream =
streamProvider.GetStream<string>(streamId);
StreamSubscriptionHandle<string> subscription =
await stream.SubscribeAsync(new MyStreamObserver());
}
IStreamProvider streamProvider =
base.GetStreamProvider("SimpleStreamProvider");
IAsyncStream<T> stream =
streamProvider.GetStream<T>(this.GetPrimaryKey(), "MyStreamNamespace");
StreamSubscriptionHandle<T> subscription =
await stream.SubscribeAsync(IAsyncObserver<T>);
Написание логики подписки
Ниже приведены рекомендации по написанию логики подписки для различных случаев: явные и неявные подписки, перематываемые и неперематываемые потоки. Основное различие между явными и неявными подписками заключается в том, что для неявных подписок зерно всегда имеет ровно одну неявную подписку на пространство имен потока. Нет способа создать несколько подписок (без кратности подписки), нет способа отмены подписки, а логика зерна должна быть присоединена только к логике обработки. Это также означает, что нет необходимости возобновлять неявную подписку. С другой стороны, для явных подписок необходимо возобновить подписку; в противном случае подписка снова приводит к тому, что зерно подписывается несколько раз.
Неявные подписки:
Для неявных подписок зерна по-прежнему необходимо подписаться на присоединение логики обработки. Это можно сделать в грануле потребителя, реализуя интерфейсы IStreamSubscriptionObserver и IAsyncObserver<T>, что позволяет грануле активироваться отдельно от подписки. Чтобы подписаться на поток, зерно создает дескриптор и вызовы await handle.ResumeAsync(this) в методе OnSubscribed(...) .
Чтобы обработать сообщения, реализуйте метод IAsyncObserver<T>.OnNextAsync(...) для получения данных потока и токена последовательности. Кроме того, метод может принимать набор делегатов, ResumeAsync представляющих методы IAsyncObserver<T> интерфейса: onNextAsync, onErrorAsyncи onCompletedAsync.
public Task OnNextAsync(string item, StreamSequenceToken? token = null)
{
_logger.LogInformation("Received an item from the stream: {Item}", item);
return Task.CompletedTask;
}
public async Task OnSubscribed(IStreamSubscriptionHandleFactory handleFactory)
{
var handle = handleFactory.Create<string>();
await handle.ResumeAsync(this);
}
public override async Task OnActivateAsync()
{
var streamProvider = this.GetStreamProvider(PROVIDER_NAME);
var stream =
streamProvider.GetStream<string>(
this.GetPrimaryKey(), "MyStreamNamespace");
await stream.SubscribeAsync(OnNextAsync);
}
Явные подписки:
Для явных подписок зерна должно вызывать SubscribeAsync подписку на поток. В результате создается подписка и подключается логика обработки. Явная подписка существует до того момента, пока зерно не отменит подписку. Если зерно деактивируется и реактивируется, оно все еще явно подписано, но логика обработки данных не присоединена. В этом случае зерна необходимо повторно подключить логику обработки. Для этого, в своем OnActivateAsync, зерно должно сначала определить свои подписки, вызвав IAsyncStream<T>.GetAllSubscriptionHandles(). Зерно должно выполнить ResumeAsync на каждом дескрипторе, с которым оно хочет продолжать обработку, или UnsubscribeAsync на любом дескрипторе, с которым закончена работа. Опционально, вы можете указать StreamSequenceToken в качестве аргумента для вызовов ResumeAsync, что приведет к тому, что эта явная подписка начнет использовать этот маркер.
public override async Task OnActivateAsync(CancellationToken cancellationToken)
{
var streamProvider = this.GetStreamProvider(PROVIDER_NAME);
var streamId = StreamId.Create("MyStreamNamespace", this.GetPrimaryKey());
var stream = streamProvider.GetStream<string>(streamId);
var subscriptionHandles = await stream.GetAllSubscriptionHandles();
foreach (var handle in subscriptionHandles)
{
await handle.ResumeAsync(this);
}
}
public async override Task OnActivateAsync()
{
var streamProvider = this.GetStreamProvider(PROVIDER_NAME);
var stream =
streamProvider.GetStream<string>(this.GetPrimaryKey(), "MyStreamNamespace");
var subscriptionHandles = await stream.GetAllSubscriptionHandles();
if (!subscriptionHandles.IsNullOrEmpty())
{
subscriptionHandles.ForEach(
async x => await x.ResumeAsync(OnNextAsync));
}
}
Потоковые маркеры порядка и последовательности
Порядок доставки событий между отдельным производителем и потребителем зависит от поставщика потоков.
С помощью SMS производитель явно управляет порядком событий, замеченных потребителем, управляя тем, как они публикуют их. По умолчанию (если SimpleMessageStreamProviderOptions.FireAndForgetDelivery параметр для поставщика SMS является false), и если производитель ожидает выполнения каждого OnNextAsync вызова, события поступают в порядке FIFO. В SMS производитель решает, как обрабатывать сбои доставки, указанные неисправным Task, возвращенным вызовом OnNextAsync.
Потоки очередей Azure не гарантируют порядок FIFO, поскольку базовые очереди Azure не гарантируют порядок в случаях сбоев (хотя при отсутствии сбоев они обеспечивают порядок FIFO). Когда производитель создает событие в очередь Azure, если операция очереди завершается ошибкой, производитель должен попытаться выполнить другую очередь и позже справиться с потенциальными повторяющимися сообщениями. На стороне доставки среда выполнения потоковой передачи извлекает событие из очереди и пытается доставить его для обработки потребителям. Среда выполнения удаляет событие из очереди только после успешной обработки. Если доставка или обработка завершается ошибкой, событие не удаляется из очереди и позже снова появляется автоматически. Среда выполнения потоковой передачи пытается снова доставить его, потенциально нарушая порядок FIFO. Это поведение соответствует нормальной семантике очередей Azure.
Определяемый приложением порядок. Для обработки описанных выше проблем с упорядочиванием приложение может дополнительно указать порядок. Для этого используется непрозрачный StreamSequenceTokenобъект, используемый для упорядочивания IComparable событий. Продюсер может передать необязательный StreamSequenceToken вызов OnNextAsync . Это StreamSequenceToken передается потребителю и доставляется вместе с событием. Таким образом, приложение может подумать и восстановить порядок независимо от среды выполнения потоковой передачи.
Перемыкаемые потоки
Некоторые потоки позволяют подписывание только начиная с последней точки во времени, а другие позволяют "вернуться в время". Эта возможность зависит от базовой технологии очередей и конкретного поставщика потоков. Например, очереди Azure позволяют обрабатывать только последние поставленные в очередь события, в то время как Центры событий позволяют воспроизводить события из произвольной точки во времени (до времени истечения срока действия). Потоки, поддерживающие возвращение назад во времени, называются перемыкаемыми потоками.
Потребитель перемотываемого потока может передать StreamSequenceTokenSubscribeAsync вызов. Среда выполнения предоставляет события, начиная с того момента StreamSequenceToken. Маркер NULL означает, что потребитель хочет получать события начиная с последней версии.
Возможность перемотки потока очень полезна в сценариях восстановления. Например, рассмотрим зерно, которое подписывается на поток и периодически выполняет контрольные точки его состояния вместе с последним маркером последовательности. При восстановлении после сбоя зерно может повторно подписаться на тот же поток из последнего маркера последовательности контрольной точки, не теряя никаких событий, созданных с момента последней контрольной точки.
Поставщик Центров событий может быть перемотан назад. Его код можно найти на сайте GitHub: Orleans/Azure/Orleans. Streaming.EventHubs. Sms (теперь широковещательный канал) и поставщики очередей Azureне перемотываются.
Автоматическое горизонтальное масштабирование без отслеживания состояния
По умолчанию Orleans целевые платформы потоковой передачи поддерживают большое количество относительно небольших потоков, каждый из которых обрабатывается одним или несколькими гранями с управлением состоянием. В совокупности обработка всех потоков сегментирована среди множества регулярных (отслеживание состояния) зерна. Код вашего приложения управляет этим сегментированием, назначая идентификаторы потоков и идентификаторы зерен и явно подписываясь на них. Цель — сегментированная обработка с отслеживанием состояния.
Однако существует также интересный сценарий автоматической горизонтальной обработки без отслеживания состояния. В этом сценарии приложение имеет небольшое количество потоков (или даже один большой поток), и цель — обработка без сохранения состояния. Примером является глобальный поток событий, в котором обработка включает декодирование каждого события и потенциально переадресация его в другие потоки для дальнейшей обработки состояния. Бесстатная масштабируемая обработка потоков может поддерживаться в Orleans с помощью StatelessWorkerAttribute grains.
Текущее состояние бестейтовой автоматически масштабируемой обработки: Это еще не реализовано. Попытка подписаться на поток из StatelessWorkerAttribute зерна приводит к неопределенному поведению. Мы рассматриваем поддержку этого варианта.
Зерна и Orleans клиенты
Orleans потоки работают равномерно между зернами и Orleans клиентами. Это означает, что вы можете использовать одни и те же API внутри грейна и в клиенте Orleans для создания и обработки событий. Это значительно упрощает логику приложения, делая специальные клиентские API, такие как Grain Observers, избыточными.
Полностью управляемая и надежная потоковая передача pub-sub
Для отслеживания подписок Orleans потоков использует компонент среды выполнения Streaming Pub-Sub, который служит точкой встречи для потребителей потоков и производителей. Pub-sub отслеживает все подписки потоков, сохраняет их и сопоставляет потребителей потоков с производителями потоков.
Приложения могут выбирать, где и как хранятся данные Pub-Sub. Сам компонент Pub-Sub реализуется как зерна (называемый PubSubRendezvousGrain), который использует Orleans декларативное сохраняемость.
PubSubRendezvousGrain использует поставщик хранилища с именем PubSubStore. Как и в случае с любым зерном, можно назначить реализацию для поставщика хранилища. Для потоковой передачи Pub-Sub можно изменить реализацию во время создания узла Silo с помощью построителя хоста Silo.
Следующая настройка настраивает Pub-Sub для хранения состояния в таблицах Azure.
var endpoint = new Uri(configuration["AZURE_TABLE_STORAGE_ENDPOINT"]!);
var credential = new DefaultAzureCredential();
hostBuilder.UseOrleans(siloBuilder =>
{
siloBuilder.AddAzureTableGrainStorage("PubSubStore",
options => options.TableServiceClient = new TableServiceClient(endpoint, credential));
});
hostBuilder.AddAzureTableGrainStorage("PubSubStore",
options => options.ConnectionString = "<Secret>");
Таким образом, данные Pub-Sub надежно хранятся в таблице Azure. Для первоначальной разработки также можно использовать хранилище памяти. Помимо Pub-Sub Orleans среда выполнения потоковой передачи предоставляет события от производителей потребителям, управляет всеми ресурсами среды выполнения, выделенными для активно используемых потоков, и прозрачно мусор собирает ресурсы среды выполнения из неиспользуемых потоков.
Настройка
Чтобы использовать потоки, необходимо включить поставщиков потоков через построители узлов или клиентов кластера silo. Пример настройки поставщика потоков:
var tableEndpoint = new Uri(configuration["AZURE_TABLE_STORAGE_ENDPOINT"]!);
var queueEndpoint = new Uri(configuration["AZURE_QUEUE_STORAGE_ENDPOINT"]!);
var credential = new DefaultAzureCredential();
hostBuilder.UseOrleans(siloBuilder =>
{
siloBuilder.AddMemoryStreams("StreamProvider")
.AddAzureQueueStreams("AzureQueueProvider",
optionsBuilder => optionsBuilder.ConfigureAzureQueue(
options => options.Configure(
opt => opt.QueueServiceClient = new QueueServiceClient(queueEndpoint, credential))))
.AddAzureTableGrainStorage("PubSubStore",
options => options.TableServiceClient = new TableServiceClient(tableEndpoint, credential));
});
hostBuilder.AddSimpleMessageStreamProvider("SMSProvider")
.AddAzureQueueStreams<AzureQueueDataAdapterV2>("AzureQueueProvider",
optionsBuilder => optionsBuilder.Configure(
options => options.ConnectionString = "<Secret>"))
.AddAzureTableGrainStorage("PubSubStore",
options => options.ConnectionString = "<Secret>");