Открыть сервис

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 формирует в буфере. Функция должна возвращать массив обработанных записей.

Процесс работы

  1. Приём данных: Источники отправляют данные в Firehose через AWS SDK, Kinesis Agent, Kinesis Data Streams, CloudWatch Logs или прямо через HTTP API.
  2. Буферизация: Данные накапливаются в буфере до достижения порога по размеру или времени.
  3. Преобразование (опционально): Если настроено, Firehose либо конвертирует формат, либо вызывает Lambda-функцию для обработки пакета.
  4. Доставка: Обработанные данные отправляются в выбранное целевое назначение. В случае ошибки (например, недоступность S3) Firehose автоматически повторяет попытку в течение заданного времени (по умолчанию 24 часа). Если повторные попытки исчерпаны, данные могут быть перенаправлены в S3-бакет для ошибок.
  5. Мониторинг: Все операции логируются в CloudWatch.

Отличия от Amazon Kinesis Data Streams

Amazon Kinesis Data Firehose часто сравнивают с другим сервисом семейства Kinesis — Amazon Kinesis Data Streams. Основные различия:

ХарактеристикаKinesis Data FirehoseKinesis 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 →