Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Потоки могут поступать в различных видах и формах. Некоторые потоки могут доставлять события через прямые tcp-каналы, а другие предоставляют события через устойчивые очереди. Различные типы потоков могут использовать различные стратегии пакетной обработки, алгоритмы кэширования или процедуры обратной обработки. Поставщики потоковых данных — это точки расширяемости в Orleans среде выполнения стриминга, которые предоставляют возможность реализовать любой тип потока, не ограничивая приложения потоковой передачи только подмножествами таких вариантов поведения. Эта точка расширяемости по своей сути аналогична поставщикам хранилища Orleans.
Поставщик потоков концентратора событий Azure
Центры событий Azure — это полностью управляемая служба приема данных в режиме реального времени, способная получать и обрабатывать миллионы событий в секунду. Он предназначен для обработки приема данных с высокой пропускной способностью, низкой задержкой из нескольких источников и последующей обработки этих данных несколькими потребителями.
Центры событий часто используются в качестве основы более крупной архитектуры обработки событий, выступающей в качестве передней двери для конвейера событий. Их можно использовать для приема данных из различных источников, включая веб-каналы социальных сетей, устройства Интернета вещей и файлы журналов. Одним из ключевых преимуществ Центров событий является возможность горизонтального масштабирования в соответствии с потребностями даже крупнейших рабочих нагрузок обработки событий. Он также высокодоступен и отказоустойчив, с несколькими репликами данных, распределенными по нескольким регионам Azure, чтобы обеспечить высокий уровень доступности.
NuGet пакет Microsoft.Orleans.Streaming.EventHubs содержит поставщик потока Event Hubs.
Поставщик потока Azure Queue (AQ)
Поставщик потоков Azure Queue (AQ) передает события через очереди Azure. На стороне производителя поставщик потоков AQ помещает события непосредственно в очередь Azure. На стороне потребителя поставщик потоков AQ управляет набором агентов извлечения , которые извлекают события из набора очередей Azure и доставляют их в код приложения, который использует их. Вы можете рассматривать агентов загрузки как распределенную "микрослужбу" — секционированный, высокодоступный и эластичный распределенный механизм. Агенты извлечения выполняются внутри одного и того же силоса, в котором размещаются приложения. Таким образом, нет необходимости запускать отдельные роли обработки Azure для извлечения данных из очередей. Среда выполнения потокового режима полностью управляет существованием агентов извлечения и их управлением, сдерживанием обратного давления, балансировкой очередей между ними и передачей очередей от вышедшего из строя агента другому агенту. Это все прозрачно для кода приложения, использующего потоки.
Пакет NuGet Microsoft.Orleans.Streaming.AzureStorage содержит поставщика потока Azure Queue Storage.
Адаптеры очередей
Различные поставщики потоков, предоставляющие события по устойчивым очередям, демонстрируют аналогичное поведение и подвергаются аналогичным реализациям. Таким образом, мы предоставляем универсальную расширяемую версию PersistentStreamProvider , которая позволяет подключать различные типы очередей без написания совершенно нового поставщика потоков с нуля.
PersistentStreamProvider использует компонент IQueueAdapter, который абстрагирует конкретные сведения о реализации очереди и предоставляет средства для добавления и удаления событий из очереди. Логика внутри PersistentStreamProvider обрабатывает все остальное. Указанный выше поставщик очередей Azure также реализуется таким образом: это экземпляр PersistentStreamProvider , который использует объект AzureQueueAdapter.
Aspire интеграция для потоковой передачи
Aspire Orleans упрощает настройку потоковой передачи, автоматически управляя подготовкой ресурсов и подключением.
Потоковая передача данных из Хранилище очередей Azure с использованием Aspire
Проект AppHost (Program.cs):
var builder = DistributedApplication.CreateBuilder(args);
var storage = builder.AddAzureStorage("storage");
var queues = storage.AddQueues("streaming");
var orleans = builder.AddOrleans("cluster")
.WithClustering(builder.AddRedis("redis"))
.WithStreaming("AzureQueueProvider", queues);
builder.AddProject<Projects.MySilo>("silo")
.WithReference(orleans)
.WithReference(queues);
builder.Build().Run();
Проект Silo (Program.cs):
var builder = Host.CreateApplicationBuilder(args);
builder.AddServiceDefaults();
builder.AddKeyedAzureQueueServiceClient("streaming");
builder.UseOrleans();
builder.Build().Run();
Подсказка
Во время локальной разработки Aspire автоматически использует эмулятор Azurite для Хранилище очередей Azure. В рабочих развертываниях Aspire подключается к реальной учетной записи служба хранилища Azure на основе конфигурации развертывания Azure.
Это важно
Необходимо вызвать AddKeyedAzureQueueServiceClient, чтобы зарегистрировать клиента очереди в контейнере внедрения зависимостей.
Orleans Поставщики потоковой передачи ищут ресурсы по имени службы с ключом, если пропустить этот шаг, Orleans не сможет разрешить клиент очереди и вызовет ошибку разрешения зависимостей во время выполнения.
Потоковая обработка данных в памяти для разработки
Для локальных сценариев разработки и тестирования можно использовать потоковую передачу в памяти, которая не требует каких-либо внешних зависимостей:
Проект AppHost (Program.cs):
var builder = DistributedApplication.CreateBuilder(args);
var orleans = builder.AddOrleans("cluster")
.WithDevelopmentClustering()
.WithMemoryStreaming("MemoryStreamProvider");
builder.AddProject<Projects.MySilo>("silo")
.WithReference(orleans);
builder.Build().Run();
Проект Silo (Program.cs):
var builder = Host.CreateApplicationBuilder(args);
builder.AddServiceDefaults();
builder.UseOrleans();
builder.Build().Run();
Предупреждение
Потоки в памяти не являются долговечными и теряются при перезапуске сайло. Используйте потоковую передачу в памяти только для разработки и тестирования— никогда не для рабочих нагрузок, требующих устойчивости сообщений.
Широковещательные каналы с Aspire
Широковещательные каналы предоставляют простой механизм публикации/подписки для трансляции сообщений всем подписчикам.
Проект AppHost (Program.cs):
var builder = DistributedApplication.CreateBuilder(args);
var orleans = builder.AddOrleans("cluster")
.WithDevelopmentClustering()
.WithBroadcastChannel("BroadcastChannel");
builder.AddProject<Projects.MySilo>("silo")
.WithReference(orleans);
builder.Build().Run();
Проект Silo (Program.cs):
var builder = Host.CreateApplicationBuilder(args);
builder.AddServiceDefaults();
builder.UseOrleans();
builder.Build().Run();
Подробную документацию по интеграции Orleans и Aspire см. в разделе Интеграция Orleans и Aspire.
Поставщик потока простых сообщений
Простой поставщик потоков сообщений, также известный как поставщик SMS, предоставляет события по протоколу TCP с помощью регулярного Orleans обмена сообщениями. Так как sms-события доставляются по ненадежным TCP-ссылкам, SMS не гарантирует надежную доставку событий и не автоматически отправляет сообщения о сбое для потоков SMS. По умолчанию вызов OnNextAsync производителя возвращает Task, представляющий состояние обработки потока потребителя. Это сообщает производителю, успешно ли потребитель получил и обработал событие. Если эта задача завершается ошибкой, производитель может снова решить отправить то же событие, чтобы обеспечить надежность на уровне приложения. Хотя доставка сообщений в потоках осуществляется в лучшем случае, SMS-потоки сами по себе надежны. То есть привязка подписчика к производителю, выполняемая Pub-Sub, полностью надежна.