Amazon Kinesis Data Firehose
Amazon Kinesis Data Firehose — это полностью управляемый сервис облачной платформы Amazon Web Services (AWS), предназначенный для потоковой загрузки (доставки) данных в реальном времени. Сервис автоматически масштабируется, не требует администрирования кластеров и позволяет с минимальной задержкой загружать данные в хранилища, аналитические сервисы и инструменты машинного обучения.
Назначение и основные функции
Основная задача Kinesis Data Firehose — упрощение процесса сбора, преобразования и доставки потоковых данных. Сервис принимает данные от различных источников (приложения, устройства IoT, логи, события), при необходимости преобразует их (например, в формат Parquet или ORC) и загружает в целевые системы.
Ключевые возможности сервиса:
- Автоматическое масштабирование: Firehose самостоятельно увеличивает пропускную способность при росте объёма входящих данных.
- Преобразование данных: возможность конвертировать данные из JSON в Parquet/ORC, а также вызывать пользовательские функции AWS Lambda для более сложных преобразований (фильтрация, обогащение, изменение формата).
- Буферизация: сервис накапливает данные в буфере перед отправкой. Буферизация может быть настроена по размеру (от 1 до 128 МБ) или по времени (от 60 до 900 секунд). Это позволяет оптимизировать затраты на запись в целевые системы.
- Мониторинг и логирование: интеграция с Amazon CloudWatch для отслеживания метрик (объём данных, задержка, ошибки) и ведения логов.
Целевые назначения (Destinations)
Amazon Kinesis Data Firehose поддерживает доставку данных в несколько типов конечных точек:
Хранилища данных
- Amazon S3 (Simple Storage Service): основной и наиболее распространённый вариант. Данные могут загружаться в S3 в виде объектов (файлов) с возможностью указания префиксов для организации каталогов (по дате, времени, ключу).
- Amazon Redshift: Firehose может загружать данные напрямую в Amazon Redshift, предварительно записывая их в промежуточный S3-бакет, а затем запуская команду COPY. Это позволяет быстро наполнять аналитические хранилища.
Аналитические сервисы
- Amazon OpenSearch Service: сервис может доставлять данные в кластеры Amazon OpenSearch Service (ранее Amazon Elasticsearch Service) для полнотекстового поиска и анализа логов.
- Amazon Elasticsearch Service (устаревшее название): поддерживается как целевое назначение для индексации данных.
Собственные конечные точки (Custom destinations)
- HTTP-эндпоинты: Firehose может отправлять данные на любой HTTP-эндпоинт, включая собственные серверы или сторонние сервисы (например, Datadog, New Relic, Splunk, MongoDB Cloud). Для этого используется специальный коннектор.
- Splunk: поддерживается прямая доставка в Splunk через HTTP Event Collector (HEC).
Преобразование и обогащение данных
Firehose предоставляет встроенные возможности для обработки данных перед загрузкой:
Форматирование
- Конвертация формата: из JSON в Apache Parquet или Apache ORC. Это значительно уменьшает объём хранимых данных и ускоряет запросы в аналитических системах (например, Amazon Athena, Amazon Redshift Spectrum).
- Сжатие: автоматическое сжатие данных (GZIP, Snappy, ZIP) перед записью в S3.
Вызов Lambda-функций
Пользователь может подключить собственную Lambda-функцию для выполнения произвольной логики преобразования:
- Фильтрация записей (удаление ненужных событий).
- Обогащение данных (добавление геолокации, идентификаторов, агрегированных значений).
- Изменение структуры (конвертация в другой JSON-формат, добавление полей).
Lambda-функция вызывается для каждого пакета данных (batch), который Firehose формирует в буфере. Функция должна возвращать массив обработанных записей.
Процесс работы
- Приём данных: Источники отправляют данные в Firehose через AWS SDK, Kinesis Agent, Kinesis Data Streams, CloudWatch Logs или прямо через HTTP API.
- Буферизация: Данные накапливаются в буфере до достижения порога по размеру или времени.
- Преобразование (опционально): Если настроено, Firehose либо конвертирует формат, либо вызывает Lambda-функцию для обработки пакета.
- Доставка: Обработанные данные отправляются в выбранное целевое назначение. В случае ошибки (например, недоступность S3) Firehose автоматически повторяет попытку в течение заданного времени (по умолчанию 24 часа). Если повторные попытки исчерпаны, данные могут быть перенаправлены в S3-бакет для ошибок.
- Мониторинг: Все операции логируются в CloudWatch.
Отличия от Amazon Kinesis Data Streams
Amazon Kinesis Data Firehose часто сравнивают с другим сервисом семейства Kinesis — Amazon Kinesis Data Streams. Основные различия:
| Характеристика | Kinesis Data Firehose | Kinesis Data Streams |
|---|---|---|
| Тип сервиса | Полностью управляемый (serverless) | Требует ручного управления шардами |
| Масштабирование | Автоматическое | Ручное (изменение числа шардов) |
| Задержка | От нескольких секунд (near real-time) | Миллисекунды (real-time) |
| Хранение данных | Нет собственного хранения (сразу доставляется) | Хранит данные до 365 дней |
| Целевые назначения | S3, Redshift, OpenSearch, HTTP, Splunk | Любые потребители (KCL, Lambda, EC2) |
| Преобразование | Встроенное (Parquet/ORC) + Lambda | Только через Lambda-потребителей |
| Стоимость | За объём загруженных данных | За шарды и объём записей |
Firehose лучше подходит для сценариев, где требуется простая, масштабируемая доставка в хранилища и аналитические сервисы без низкой задержки. Data Streams — для случаев, когда необходима задержка в миллисекундах, возможность повторного чтения данных или реализация сложной потоковой обработки (например, с помощью Apache Flink).
Варианты использования
- Сбор логов веб-приложений: логи с веб-серверов (Nginx, Apache) или приложений отправляются в Firehose, который преобразует их в Parquet и загружает в S3 для последующего анализа с помощью Amazon Athena.
- Аналитика IoT: данные с датчиков и устройств собираются, фильтруются (например, удаляются невалидные показания) и загружаются в Amazon Redshift для построения дашбордов.
- Мониторинг безопасности: журналы безопасности (CloudTrail, VPC Flow Logs) доставляются в Amazon OpenSearch Service для поиска и визуализации.
- Резервное копирование потоковых данных: все события из Kinesis Data Streams могут быть автоматически скопированы в S3 через Firehose для долговременного хранения.
Ограничения
- Задержка: минимальная задержка доставки составляет несколько секунд, что не подходит для задач, требующих реакции в реальном времени (миллисекунды).
- Преобразование: встроенные преобразования ограничены конвертацией формата и вызовом Lambda. Для сложной обработки (агрегации, оконные функции) требуется использование Kinesis Data Analytics или Data Streams.
- Размер записи: максимальный размер одной записи — 1 МБ (до преобразования). После преобразования размер может увеличиться, но не более 4 МБ.
- Пропускная способность: хотя сервис масштабируется автоматически, существуют региональные лимиты по умолчанию (например, 5000 записей/сек для одного потока Firehose), которые можно увеличить через запрос в службу поддержки AWS.
Стоимость
Ценообразование Kinesis Data Firehose основано на объёме загруженных данных (в гигабайтах) и количестве вызовов Lambda-функций для преобразования. Дополнительно взимается плата за хранение данных в S3, Redshift или OpenSearch. Сжатие и конвертация в Parquet/ORC выполняются без дополнительной платы.
Источники
- AWS Documentation: Amazon Kinesis Data Firehose Developer Guide
- AWS re:Invent 2023: «Best practices for Amazon Kinesis Data Firehose» (презентация)
- AWS Whitepaper: «Streaming Data Solutions on AWS»
- Документация Amazon Web Services на русском языке (раздел «Amazon Kinesis Data Firehose»)
BFOmetr — база данных и аналитика по компаниям России.
На главную BFOmetr →