Orleans Краткое руководство по потоковой передаче

В этом руководстве показано, как быстро настроить и использовать Orleans Streams. Дополнительные сведения о функциях потоковой передачи см. в других частях этой документации.

Обязательные конфигурации

В этом руководстве вы будете использовать поток, использующий память и грейновое обмен сообщений, для передачи потоковых данных подписчикам. Поставщик хранилища в памяти используется для хранения списков подписок. Использование механизмов на основе памяти для потоковой передачи и хранения предназначено только для локальной разработки и тестирования, а не для рабочих сред.

Orleans для потоковой передачи требуется пакет Microsoft.Orleans.Streaming NuGet. Этот пакет предоставляет функции потоковой передачи для клиента и сервера, включая AddMemoryStreams метод расширения, используемый в этом руководстве.

На силосе, где silo является ISiloBuilder, вызовите AddMemoryStreams:

silo.AddMemoryStreams("StreamProvider")
    .AddMemoryGrainStorage("PubSubStore");

На клиенте кластера, где client — это IClientBuilder, вызовите AddMemoryStreams.

client.AddMemoryStreams("StreamProvider");

В этом руководстве используйте простой поток на основе сообщений, который использует обмен сообщениями с зерном для отправки потоковых данных подписчикам. Используйте поставщик хранилища в памяти для хранения списков подписок; это не мудрый выбор для реальных рабочих приложений.

На силосе, где hostBuilder является ISiloHostBuilder, вызовите AddSimpleMessageStreamProvider:

hostBuilder.AddSimpleMessageStreamProvider("SMSProvider")
           .AddMemoryGrainStorage("PubSubStore");

На клиенте кластера, где clientBuilder — это IClientBuilder, вызовите AddSimpleMessageStreamProvider.

clientBuilder.AddSimpleMessageStreamProvider("SMSProvider");

Примечание.

По умолчанию сообщения, передаваемые по простому потоку сообщений, считаются неизменяемыми и могут передаваться по ссылке на другие зерна. Чтобы отключить это поведение, настройте поставщика SMS для отключения SimpleMessageStreamProviderOptions.OptimizeForImmutableData.

siloBuilder
    .AddSimpleMessageStreamProvider(
        "SMSProvider",
        options => options.OptimizeForImmutableData = false);

Вы можете создавать потоки, отправлять данные с помощью их в качестве производителей и получать данные в качестве подписчиков.

Создание событий

Создавать события для потоков довольно легко. Сначала получите доступ к поставщику потоков, определенному в конфигурации ранее ("StreamProvider"), а затем выберите поток и отправьте в него данные.

// Pick a GUID for a chat room grain and chat room stream
var guid = new Guid("some guid identifying the chat room");
// Get one of the providers which we defined in our config
var streamProvider = GetStreamProvider("StreamProvider");
// Get the reference to a stream
var streamId = StreamId.Create("RANDOMDATA", guid);
var stream = streamProvider.GetStream<int>(streamId);

Создавать события для потоков довольно легко. Сначала получите доступ к поставщику потоков, определенному в конфигурации ранее ("SMSProvider"), а затем выберите поток и отправьте в него данные.

// Pick a GUID for a chat room grain and chat room stream
var guid = new Guid("some guid identifying the chat room");
// Get one of the providers which we defined in our config
var streamProvider = GetStreamProvider("SMSProvider");
// Get the reference to a stream
var stream = streamProvider.GetStream<int>(guid, "RANDOMDATA");

Как видно, в потоке есть GUID и пространство имен. Это упрощает идентификацию уникальных потоков. Например, пространство имен для комнаты чата может быть "Комнаты", а GUID может быть GUID владельца RoomGrain.

Здесь используйте GUID известной комнаты чата. Используя метод потока OnNextAsync, отправьте данные в него. Давайте сделаем это внутри таймера с помощью случайных чисел. Вы также можете использовать любой другой тип данных для потока.

RegisterTimer(_ =>
{
    return stream.OnNextAsync(Random.Shared.Next());
},
null,
TimeSpan.FromMilliseconds(1_000),
TimeSpan.FromMilliseconds(1_000));

Подписывайтесь и получайте потоковые данные

Для получения данных можно использовать неявные и явные подписки, описанные более подробно в разделе "Явные и неявные подписки". В этом примере используются неявные подписки, которые проще. Если тип зерна хочет неявно подписаться на поток, он использует атрибут [ImplicitStreamSubscription(namespace)].

В вашем случае определите ReceiverGrain следующим образом:

[ImplicitStreamSubscription("RANDOMDATA")]
public class ReceiverGrain : Grain, IRandomReceiver

При отправке данных в потоки в пространстве имен RANDOMDATA (как в примере с таймером), зерно типа ReceiverGrain с тем же Guid, что и поток, получает сообщение. Даже если активации зерна в настоящее время отсутствуют, среда выполнения автоматически создает новый и отправляет в него сообщение.

Для этого выполните процесс подписки, задав OnNextAsync метод получения данных. Для этого ReceiverGrain должно вызвать что-то подобное в OnActivateAsync:

// Create a GUID based on our GUID as a grain
var guid = this.GetPrimaryKey();

// Get one of the providers which we defined in config
var streamProvider = GetStreamProvider("StreamProvider");

// Get the reference to a stream
var streamId = StreamId.Create("RANDOMDATA", guid);
var stream = streamProvider.GetStream<int>(streamId);

// Set our OnNext method to the lambda which simply prints the data.
// This doesn't make new subscriptions, because we are using implicit
// subscriptions via [ImplicitStreamSubscription].
await stream.SubscribeAsync<int>(
    async (data, token) =>
    {
        Console.WriteLine(data);
        await Task.CompletedTask;
    });
// Create a GUID based on our GUID as a grain
var guid = this.GetPrimaryKey();

// Get one of the providers which we defined in config
var streamProvider = GetStreamProvider("SMSProvider");

// Get the reference to a stream
var stream = streamProvider.GetStream<int>(guid, "RANDOMDATA");

// Set our OnNext method to the lambda which simply prints the data.
// This doesn't make new subscriptions, because we are using implicit
// subscriptions via [ImplicitStreamSubscription].
await stream.SubscribeAsync<int>(
    async (data, token) =>
    {
        Console.WriteLine(data);
        await Task.CompletedTask;
    });

Все готово! Теперь единственное требование состоит в том, чтобы что-то инициировало создание зерна-производителя. Затем он регистрирует таймер и начинает отправлять случайные целые числа всем заинтересованным сторонам.

Опять же, в этом руководстве пропускается множество сведений и предоставляется только общий обзор. Ознакомьтесь с другими частями этого руководства и другими ресурсами по Rx, чтобы получить хорошее представление о том, что доступно и как это работает.

Реактивное программирование может быть мощным подходом к решению многих проблем. Например, можно использовать LINQ в подписчике для фильтрации номеров и выполнения различных интересных операций.

См. также

Orleans API программирования потоков