Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Широкий спектр технологий уже существует для создания систем потоковой обработки. К ним относятся системы для хранения потоковых данных (например, Центров событий и Kafka) и систем для экспресс-операций вычислений через потоковые данные (например, Azure Stream Analytics, Apache Storm и Apache Spark Streaming). Это отличные системы, которые позволяют создавать эффективные конвейеры обработки потока данных.
Ограничения существующих систем
Однако эти системы не подходят для точного вычисления свободной формы через потоковые данные. Вычислительные системы потоковой передачи, упомянутые выше, позволяют указать единый граф потока данных операций, применяемых одинаково ко всем элементам потока. Это мощная модель, когда данные являются универсальными, и вы хотите выразить тот же набор операций преобразования, фильтрации или агрегирования по этим данным. Но другие варианты использования требуют выражения принципиально разных операций над различными элементами данных. В некоторых из этих случаев в процессе обработки иногда может потребоваться выполнить внешний вызов, например вызов произвольного REST API. Единые подсистемы обработки потоков данных либо не поддерживают эти сценарии, поддерживают их в ограниченном и ограниченном порядке, либо неэффективны в их поддержке. Это связано с тем, что они по сути оптимизированы для большого объема аналогичных элементов и обычно ограничены с точки зрения экспрессивности и обработки. Orleans Потоки нацелены на эти другие сценарии.
Мотивация
Все началось с запросов от Orleans пользователей с целью поддержки возврата последовательности элементов из вызова метода грейна. Как вы можете себе представить, это была только вершина айсберга; им нужно было гораздо больше.
Типичный сценарий для Orleans Streams — это когда у вас есть потоки для каждого пользователя и требуется выполнять разные обработки для каждого пользователя в контексте этого отдельного пользователя. У вас могут быть миллионы пользователей, но некоторые заинтересованы в погоде и подписываются на оповещения о погоде для определенного места, а другие заинтересованы в спортивных мероприятиях; кто-то другой может отслеживать состояние определенного полета. Для обработки этих событий требуется другая логика, но вы не хотите запускать два независимых экземпляра потоковой обработки. Некоторые пользователи могут быть заинтересованы только в определенной ценной бумаге и только в том случае, если применяется определенное внешнее условие— условие, которое, возможно, не является частью потоковых данных (и поэтому требует динамической проверки во время выполнения обработки).
Пользователи изменяют свои интересы все время, поэтому их подписки на определенные потоки событий приходят и идут динамически. Таким образом, топология потоковой передачи меняется динамически и быстро. Поверх этого логика обработки на пользователя развивается и динамически изменяется на основе состояния пользователя и внешних событий. Внешние события могут изменить логику обработки для конкретного пользователя. Например, в системе обнаружения обмана игр при обнаружении нового метода обмана логика обработки требует обновления с помощью нового правила для обнаружения этого нарушения. Это необходимо сделать, конечно, без нарушения текущего конвейера обработки. Подсистемы обработки потоков потоков массовых данных не были созданы для поддержки таких сценариев.
Это почти не говорит о том, что такая система должна работать на нескольких сетевых компьютерах, а не только на одном узле. Таким образом, логика обработки должна быть распределена масштабируемо и эластично по кластеру серверов.
Новые требования
Четыре основных требования были определены для системы потоковой обработки для целей приведенных выше сценариев:
- Гибкая логика потоковой обработки
- Поддержка высокодинамовых топологий
- Тонкая степень детализации потока
- Распределение
Гибкая логика потоковой обработки
Система должна поддерживать различные способы выражения логики потоковой обработки. Существующие системы, упомянутые выше, требуют, чтобы разработчики писали декларативный граф вычислений потока данных, обычно следуя функциональному стилю программирования. Это ограничивает экспрессивность и гибкость логики обработки. Orleans потоки равнодушны к тому, как выражается логика обработки. Его можно выразить как поток данных (например, с помощью реактивных расширений (Rx) в .NET), функциональной программы, декларативного запроса или общей императивной логики. Логика может быть с сохранением состояния или без сохранения состояния, может иметь побочные эффекты и может вызывать внешние действия. Вся власть переходит к разработчику.
Поддержка динамических топологий
Система должна позволять динамическое развитие топологий. Существующие системы, упомянутые выше, обычно ограничены статическими топологиями, фиксированными во время развертывания, которые не могут развиваться во время выполнения. В следующем примере выражения потока данных все хорошо и просто, пока не потребуется изменить его:
Stream.GroupBy(x=> x.key).Extract(x=>x.field).Select(x=>x+2).AverageWindow(x, 5sec).Where(x=>x > 0.8) *
Измените пороговое условие в Where фильтре, добавьте инструкцию или добавьте Select другую ветвь в граф потока данных и создайте новый выходной поток. В существующих системах это невозможно без разрыва всей топологии и перезапуска потока данных с нуля. Практически эти системы выполняют контрольные точки существующего вычисления и могут перезапуститься с последней контрольной точки. Тем не менее, такой перезапуск является разрушительным и дорогостоящим для веб-службы, производящей результаты в режиме реального времени. Такой перезапуск становится особенно непрактичным при работе с большим количеством таких выражений, выполняемых с аналогичными, но разными параметрами (на пользователя, на устройство и т. д.), которые постоянно изменяются.
Система должна разрешить развитие графа потоковой обработки во время выполнения путем добавления новых ссылок или узлов в граф вычислений или изменения логики обработки в вычислительных узлах.
Тонкая степень детализации потока
В существующих системах наименьшая единица абстракции обычно представляет собой весь поток (топология). Однако во многих целевых сценариях требуется, чтобы отдельный узел или ссылка в топологии сами по себе были логической сущностью. Таким образом, каждая сущность может управляться независимо. Например, в топологии большого потока, состоящей из нескольких соединений, разные соединения могут иметь разные характеристики и осуществляться через различные физические носители. Некоторые ссылки могут переходить по сокетам TCP, а другие используют надежные очереди. Различные ссылки могут иметь разные гарантии доставки. Разные узлы могут иметь различные стратегии контрольных точек, и их логика обработки может быть выражена в разных моделях или даже на разных языках. Такая гибкость обычно невозможна в существующих системах.
Единица абстракции и гибкости аналогична сравнению СОА (сервис-ориентированных архитектур) и акторов. Системы субъектов обеспечивают большую гибкость, так как каждый субъект, по сути, является независимо управляемым "крошечной службой". Аналогичным образом, система потоковой передачи должна обеспечить такой точный контроль.
Распределение
И, конечно, система должна иметь все свойства "хорошей распределенной системы". Это включает в себя:
- Масштабируемость: поддерживает большое количество потоков и вычислительных элементов.
- Эластичность: Позволяет добавлять или удалять ресурсы для масштабирования в зависимости от нагрузки.
- Надежность: устойчивость к сбоям.
- Эффективность: эффективно использует базовые ресурсы.
- Скорость реагирования. Включает сценарии почти в режиме реального времени.
Это были требования к созданию Orleans потоковой передачи.
Уточнение: Orleans в настоящее время не поддерживает прямое написание декларативных выражений потока данных, как показано в приведенном выше примере. Текущие Orleans API потоковой передачи являются более низкоуровневыми строительными блоками, как описано в Orleans API потоковой передачи.