Разработка трубопроводов обработки данных в реальном времени с использованием технологий без серверов
Организации, которые полагаются на пакетную обработку, часто реагируют на часы или даже дни после событий. Напротив, конвейеры обработки данных в реальном времени позволяют немедленно принимать решения, обнаруживать аномалии и персонализировать пользовательский опыт. Безсерверные технологии устраняют операционные накладные расходы на управление серверами, что позволяет строить эти трубопроводы с минимальным бременем инфраструктуры. Объединив вычисления, управляемые событиями, управляемый поток поглощения и масштабируемое хранилище, команды могут создавать производственные системы реального времени, которые автоматически масштабируются от нуля до тысяч событий в секунду.
Что такое бессерверные технологии?
Бессерверные вычисления — это модель облачного исполнения, где облачный провайдер динамически управляет распределением и предоставлением серверов. Разработчики пишут и развертывают код в виде функций или контейнеров, а провайдер обрабатывает масштабирование, патчи и доступность. Термин «безсерверный» не означает, что серверы отсутствуют; скорее, управление сервером абстрагировано. Крупные провайдеры, такие как AWS Lambda, Azure Functions и Google Cloud Functions являются наиболее распространенными вычислительными службами. Они выполняют код в ответ на события — например, HTTP-запрос, новый файл в хранилище или сообщение, поступающее в очередь. Биллинг основан на продолжительности выполнения и потреблении ресурсов, а не на нерабочем потенциале. Это делает бессерверный особенно привлекательным для переменных рабочих нагрузок и обработки данных в режиме реального времени, где объем данных может непредсказуемо увеличиваться.
Помимо вычислений, бессерверные включают управляемые сервисы для приема данных, хранения, обмена сообщениями и аналитики — все это может быть собрано в конвейер без предоставления одной виртуальной машины. Ключевые характеристики включают автоматическое масштабирование, ценообразование с оплатой за использование и встроенную отказоустойчивость. При построении трубопроводов в реальном времени эти черты приводят к более низкой задержке и снижению операционной сложности по сравнению с традиционными серверными архитектурами.
Ключевые компоненты трубопроводов данных в реальном времени
Трубопровод данных в реальном времени представляет собой непрерывный поток, в котором данные поступают, обрабатываются, хранятся и действуют в течение нескольких секунд или миллисекунд. Основные строительные блоки остаются согласованными на облачных платформах:
- Потребление данных — точка входа, которая захватывает события от производителей (датчики IoT, мобильные приложения, журналы веб-серверов, базы данных). Управляемые потоковые сервисы, такие как Amazon Kinesis Data Streams, Azure Event Hubs и Google Cloud Pub/Sub, предназначены для обработки высокопроизводительных, прочных событий. Они буферизируют события и делают их доступными для потребителей в порядке.
- Обработка данных — преобразование, фильтрация, агрегация, обогащение или анализ событий по мере их прохождения по конвейеру. Функции без сервера — AWS Lambda, Azure Functions, Google Cloud Functions — являются наиболее легким вариантом для обработки без состояния, событийных функций. Для более сложных преобразований или государственных операций (например, оконные агрегации) провайдеры предлагают движки обработки потоков без сервера, такие как AWS Kinesis Data Analytics для Apache Flink, Azure Stream Analytics или Google Cloud Dataflow (который также работает на безсерверной модели с автомасштабированием).
- Хранение данных — место назначения, где обработанные результаты сохраняются для аналитики, панелей инструментов или долгосрочного хранения. Варианты варьируются от магазинов с ключевыми значениями (Amazon DynamoDB, Azure Cosmos DB) до колоночных баз данных (Google BigQuery, Amazon Redshift Serverless) и магазинов объектов (Amazon S3, Azure Blob Storage). Выбор зависит от шаблонов запросов, требований к задержке и стоимости.
- Визуализация и мониторинг — инструменты, обеспечивающие приборные панели в реальном времени, оповещение и наблюдаемость. Управляемые BI-сервисы, такие как Amazon QuickSight, Microsoft Power BI (подключенные через потоковые наборы данных) и Google Looker Studio, могут потреблять живые данные. Кроме того, мониторинг самого конвейера имеет решающее значение: такие сервисы, как Amazon CloudWatch, Azure Monitor и Google Cloud Operations Suite, вызовы функций отслеживания, задержка потока, частота ошибок и пропускная способность.
Эти компоненты должны быть соединены вместе с обменом сообщениями, безопасностью и оркестровкой. Безсерверные технологии делают каждую часть независимо масштабируемой, и клей часто обеспечивается слоем интеграции событий облачной платформы.
Архитектурные шаблоны для бессерверных трубопроводов реального времени
Хотя строительные блоки являются общими, выбранная вами архитектура зависит от характера данных и требуемых гарантий. Доминируют три шаблона:
Вентилятор с очередями сообщений
События достигают одной точки приема (например, концентратора событий или потока) и затем раздуваются до нескольких бессерверных функций или стоков хранения. Этот шаблон идеален, когда одно и то же сырое событие должно вызывать несколько независимых действий - например, обновление панели мониторинга в реальном времени, запись в холодное хранилище и отправка оповещения. Использование отдельных функций Lambda или функций Azure, которые каждый подписывается на один и тот же поток или очередь, позволяет независимое масштабирование и избежать связи. Недостатком является потенциальная дублирующая обработка или упорядочение сложностей, если функции не являются идемпотентными.
Цепная обработка с шагами функций
Некоторые трубопроводы требуют последовательных этапов обработки, где выход одной функции подается в следующую. Вместо того, чтобы организовывать эти вызовы вручную с помощью кода, сервисные оркестрторы, такие как AWS Step Functions, Azure Logic Apps или Google Cloud Workflows, координируют последовательность бессерверных функций. Это полезно для ETL-подобных преобразований, где данные должны быть проверены, обогащены, а затем агрегированы. Оркестр управляет повторными запросами, обработкой ошибок и параллельными ветвями, упрощая общую логику трубопровода. Задержка в реальном времени выше, чем прямое вызов функции к функции, но компромисс лучше наблюдаемость и устойчивость.
Потоковая обработка с использованием Stateful Compute
Для случаев использования, которые включают в себя оконные агрегаты (например, подсчет кликов в минуту) или сложную обработку событий (совпадение шаблонов между событиями), функции без состояния недостаточны. Двигатели обработки потоков без сервера, такие как Apache Flink на Kinesis Data Analytics или Google Dataflow, обрабатывают состояние, временные окна и точно один раз семантику. Эти службы работают без сервера — вы определяете логику обработки (SQL или Java / Python) и работники автомасштаба платформы. Этот шаблон является самым мощным для аналитики в реальном времени, но требует тщательного управления размером состояния и контрольных точек, чтобы избежать затрат.
Строительство трубопровода: пример AWS
Чтобы обосновать концепции, рассмотрим конкретный сценарий: прием данных веб-потока кликов, обработка их для подсчета просмотров страниц по URL в одноминутных окнах и хранение результатов для панели мониторинга в режиме реального времени. Используя полностью бессерверные сервисы AWS:
- Поток данных: Поток данных Kinesis с двумя осколками (масштабами по мере необходимости). Каждый осколок может проглотить 1 МБ/с или 1000 записей/с. Производители — такие как веб-приложение или журналирование CloudFront — отправляют события JSON в поток.
- Обработка данных: Функция Lambda запускается потоком Kinesis (с использованием отображения источника событий). Функция считывает партии записей, анализирует JSON и подсчитывает поле «url». Однако функции Lambda не имеют состояния, и каждый вызов обрабатывает микро-пакет. Для выполнения оконного подсчета можно записать счета в таблицу DynamoDB с TTL, затем использовать другую Lambda для агрегирования. Альтернативно, использовать Kinesis Data Analytics для Flink с приложением SQL, которое запускает «SELECT url, COUNT(*) FROM my stream GROUP BY url, TUMBLE(event time, INTERVAL '1' MINUTE)». Приложение Flink выводит результаты в поток вывода Kinesis Data Analytics.
- Хранение: Выходной поток запускает другую функцию Lambda, которая записывает агрегированные подсчеты (URL, count, время окончания окна) в DynamoDB с TTL, скажем, 24 часа. Одновременно сырые события могут быть архивированы в S3 с использованием Kinesis Firehose для последующего анализа.
- Визуализация: Amazon QuickSight подключается к DynamoDB через Athena (с помощью разъема Athena DynamoDB) для создания панели инструментов в реальном времени, которая обновляется каждую минуту. Альтернативно, используйте пользовательское приложение с API Serverless WebSocket для продвижения обновлений для клиентов браузера.
Весь этот трубопровод не использует экземпляры EC2, не имеет ручного масштабирования и несет расходы только при потоках данных. Функции Lambda, пропускная способность чтения / записи DynamoDB и часы Kinesis являются основными драйверами затрат. Мониторинг обрабатывается приборными панелями CloudWatch и сигнализациями о возрасте потока (millisBehindLatest) для обнаружения замедлений.
Преимущества использования бессерверных трубопроводов в реальном времени
- Истинная эластичность: Безсерверные сервисы масштабируются от нуля до тысяч одновременных исполнений за секунды.Во время флеш-продажи или вирусного события автоматические разделы трубопровода работают в большем количестве экземпляров функций или осколков потока — планирование пропускной способности не требуется.
- Стоимость-эффективность: Платите только за потребленные ресурсы. Функции оплачиваются за миллисекунду исполнения; потоковое хранилище — за GB-час; операции с базами данных — за чтение/запись. Для неработающей инфраструктуры нет затрат. Для пиковых рабочих нагрузок бессерверные могут быть на 70% дешевле, чем обеспеченные серверы.
- Сокращение операционных накладных расходов: Никаких исправлений сервера, никаких обновлений ОС, никакого прогнозирования пропускной способности. Команда может сосредоточиться на бизнес-логике и качестве данных, а не на управлении инфраструктурой.
- Гибкость и интеграция: Каждый облачный провайдер предлагает десятки источников событий, которые могут запускать функции или потоковые процессоры — потоки изменения базы данных (DynamoDB Streams, Change Data Capture from RDS), загрузки файлов (S3 Events), веб-хуки и многое другое.
- Изоляция ошибок: Сбой в одной функции вызова не приводит к сбоям в других частях трубопровода. Такие службы, как Lambda имеют встроенную логику повторных попыток и DLQ (очереди мертвой буквы). Постоянные потоковые процессоры могут проверять точки и восстанавливаться после сбоев без потери данных.
Проблемы и соображения
Бессерверные трубопроводы в реальном времени являются мощными, но создают конкретные проблемы, которые архитекторы должны решать:
- Холодные старты: Когда бессерверная функция не запускается в течение периода, платформа должна инициализировать новый контейнер, добавляя задержку (часто 100-500 мс. Для трубопроводов реального времени, где задержка до 100 мс критическая, холодные запуски могут быть проблематичными.Смягчения включают обеспеченную параллель (сохранение установленного количества теплых экземпляров), использование более легких сред выполнения (например, Node.js против Java) или использование служб обработки потоков, которые всегда теплые.
- Управление состоянием: Функции не имеют состояния по дизайну. Если трубопроводу необходимо соотносить события во времени (например, обнаруживать сеанс пользователя), состояние должно храниться внешне (DynamoDB, ElastiCache или процессор потоков без сервера). Это добавляет задержку и стоимость. Выбор правильного государственного хранилища и управление TTL являются необходимыми.
- Точно-разовые гарантии: Достижение точной однократной обработки в бессерверных трубопроводах затруднено. Функции Lambda, вызванные потоком, могут получать дублирующие записи из-за повторных запросов. Идемпотентная обработка (например, использование уникальных идентификаторов событий и восходящий к хранению) является обязательной. Двигатели обработки потока, такие как Flink, могут обеспечить точно один раз семантику в трубопроводе к поглотителям нисходящего потока, но сами поглотители также должны поддерживать его.
- Мониторинг и отладка:] При многих эфемерных вызовах функций традиционный анализ журналов становится подавляющим. Необходимы централизованные журналы (CloudWatch Logs, Azure Log Analytics), распределенное отслеживание (AWS X-Ray, OpenTelemetry) и структурированные журналы. Тревоги должны быть установлены на метриках здоровья трубопровода, а не только на функциональных ошибках.
- Vendor Lock-in: Каждый облачный провайдер имеет свой собственный вкус бессерверных сервисов и интеграции событий. Поток, построенный на Kinesis + Lambda + DynamoDB, не является непосредственно переносимым на Azure Event Hubs + Azure Functions + Cosmos DB. Mitigate путем абстрагирования логики трубопровода в переносной код (например, с использованием стандарта CloudEvents) и с использованием фреймворков обработки потоков с открытым исходным кодом, таких как Apache Flink или Apache Kafka.
Стратегии оптимизации затрат
Бессерверные модели ценообразования требуют тщательного проектирования, чтобы избежать сюрпризов:
- Батч События: Функции могут обрабатывать несколько записей на вызов. С помощью Kinesis настройте размер партии и окно партии, чтобы минимизировать количество вызовов. Например, обработка 1000 записей в одном исполнении функций стоит столько же, сколько и в одном исполнении — намного дешевле, чем 1000 отдельных вызовов.
- Вычисление правильного размера: Выделение памяти Lambda напрямую коррелирует с процессором и пропускной способностью сети. Для преобразований данных, связанных с процессором (например, JSON-парсинг, сжатие), увеличение памяти (и, следовательно, CPU) может сократить время выполнения и снизить общую стоимость (поскольку стоимость = память * продолжительность). Функции профиля с AWS Lambda Power Tuning для поиска оптимальной настройки памяти.
- Использовать управляемые потоковые процессоры для большого объема:] Для пропускной способности выше нескольких тысяч записей в секунду Lambda может стать дорогой из-за платы за запрос. Kinesis Data Analytics или Azure Stream Analytics, имея базовую почасовую стоимость, часто оказываются дешевле на миллион событий, потому что они пакетно обрабатывают внутри и заряжают на потоковый блок.
- Сжатие данных: Сжатие событий перед отправкой в поток снижает затраты на хранение и время выполнения Lambda. Gzip или snappy могут значительно уменьшить размер полезной нагрузки.
- Переносные TTL: Временное хранение (Политики жизненного цикла DynamoDB, S3) должно иметь автоматический срок годности. Обработанные промежуточные результаты, которые не нужны после того, как окно может быть выброшено.
Рассмотрение вопросов безопасности
В режиме реального времени конвейеры часто обрабатывают конфиденциальные данные. Лучшие практики безопасности без сервера включают:
- Наименее привилегированный IAM: Каждая функция должна иметь узкую роль IAM, которая предоставляет только необходимые действия на конкретных ресурсах. Например, функция Lambda, читающая от Kinesis, должна иметь «GetRecords», «DescribeStream» и «ListShards» на этом конкретном потоке, не более того. Используйте ключи состояния, чтобы ограничиться конкретными исходными конечными точками VPC, если это необходимо.
- Шифровать данные в транзите и в режиме покоя: Включить шифрование на потоках Kinesis (AWS KMS), таблицах DynamoDB и вёдрах S3. Используйте TLS для любых внешних вызовов API. Функции без сервера также могут использовать переменные среды с шифрованием KMS для секретов.
- Размещение VPC: Если трубопроводу необходимо получить доступ к ресурсам внутри VPC (например, частной базы данных), разместите функции Lambda в VPC с соответствующими группами безопасности и подсетями. Имейте в виду, что функции VPC Lambda имеют более длительные холодные запуски и требуют шлюза NAT для доступа в Интернет — что добавляет стоимость.
- Вводная валидация и санизация: Поскольку события могут происходить из ненадежных источников, бессерверные функции должны проверять и дезинфицировать все входы для предотвращения атак инъекции или искаженных данных от сбоя трубопровода. Используйте библиотеки проверки схемы (например, JSON Schema) в точке приема внутрь.
Реальные случаи использования
Бессерверные трубопроводы реального времени развернуты в различных отраслях промышленности:
- Персонализация электронной коммерции: Потоковые данные для обновления моделей рекомендаций в режиме реального времени. Функции Lambda обогащают события профилями пользователей из DynamoDB, затем нажимают на кэш, такой как ElastiCache для движка рекомендаций. Результаты отображаются на веб-сайте в течение нескольких секунд.
- Аномалия обнаружения IoT: Устройства отправляют телеметрию (температура, вибрация) в Azure Event Hubs. Функция без сервера в Azure Functions запускает облегченную модель обнаружения аномалий (например, с использованием ML.NET или Python scikit-learn) и запускает оповещение через Azure Logic Apps, если значения превышают пороговые значения. Обработанные данные хранятся в Time Series Insights.
- Обнаружение финансового мошенничества: События транзакций проходят через Google Cloud Pub/Sub к облачным функциям, а затем к Bigtable. Работа по обработке потока с использованием Dataflow (Apache Beam) применяет сопоставление оконных шаблонов для обнаружения тестирования карт или попыток захвата аккаунта. Подозреваемые транзакции помечаются и отправляются в систему «человек в петле».
- Аналитика в масштабе: Журналы приложений поступают через Kinesis Firehose непосредственно в S3 и Elasticsearch (Amazon OpenSearch Serverless). Функции Lambda анализируют и структурируют журналы перед индексацией. Панели управления в OpenSearch Dashboards обеспечивают частоту ошибок в реальном времени и процентили задержки.
Внешние ресурсы
Для более глубоких погружений обратитесь к этим официальным документам и руководствам:
- AWS: Руководство для разработчиков потоков данных Amazon Kinesis
- Azure: Введение в Azure Stream Analytics
- Google Cloud: Проводные трубопроводы с потоком данных
- Серверная среда: Центр обучения без сервера
Заключение
Безсерверные технологии созрели для поддержки требовательных конвейеров обработки данных в реальном времени. Используя управляемые службы приема, вычисления, основанные на событиях, и масштабируемое хранилище, команды могут создавать системы, которые реагируют на данные в течение нескольких секунд, минимизируя работу инфраструктуры. Ключ заключается в выборе правильного шаблона - безгосударственные функции для простых преобразований, управляемые потоковые процессоры для государственных оконных аналитики и оркестроры для многоступенчатых рабочих процессов. С тщательным вниманием к холодным запускам, управлению состоянием и мониторингу затрат, бессерверные трубопроводы в реальном времени могут обеспечить эластичность и экономичность, которые требуют современные приложения. Облачные провайдеры продолжают инвестировать в более низкую задержку, лучшую обработку состояния и упрощенную интеграцию, делая этот подход все более жизнеспособным для критически важных потоковых рабочих нагрузок.