Потоки Kinesis, Firehose и аналитика в реальном времени
Получайте потоковые данные с помощью Kinesis Data Streams, передавайте их в S3 или Redshift через Firehose и анализируйте их в реальном времени с помощью Managed Service for Apache Flink.
«Потоки Kinesis, Firehose и аналитика в реальном времени» — бесплатный урок AWS Solutions Architect на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AWS Solutions Architect, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AWS Solutions Architect содержит 4 уроков всего.
Обзор семейства Kinesis
Amazon Kinesis — это семейство служб для сбора, обработки и анализа потоковых данных в реальном времени. Три основные службы: Kinesis Data Streams (малая задержка и пользовательская обработка), Kinesis Data Firehose (полностью управляемая доставка в S3/Redshift/OpenSearch) и Managed Service for Apache Flink (ранее Kinesis Data Analytics) для выполнения SQL и обработки с помощью Flink в реальном времени. Каждая служба предназначена для своей задачи в потоковом конвейере.
Архитектура Kinesis Data Streams
Поток данных Kinesis — это надёжный упорядоченный журнал, разделённый на сегменты. Каждый сегмент обеспечивает пропускную способность записи 1 МБ/с и чтения 2 МБ/с. По умолчанию записи данных хранятся 24 часа (срок можно увеличить до 7 или 365 дней). Производители записывают данные в поток, а потребители — Lambda, приложения KCL, Firehose или Flink — параллельно считывают их из одного или нескольких сегментов. После записи записи неизменяемы.
# Create a Kinesis Data Stream with 4 shards
aws kinesis create-stream \
--stream-name clickstream \
--shard-count 4
# Put a record into the stream
aws kinesis put-record \
--stream-name clickstream \
--partition-key 'user-123' \
--data 'eyJldmVudCI6ICJjbGljayJ9'Сегменты, пропускная способность и масштабирование
Количество сегментов определяет общую пропускную способность потока. Вы можете разделить сегмент, чтобы удвоить пропускную способность, или объединить два сегмента, чтобы снизить расходы. Используйте расширенную раздачу, чтобы предоставить каждому зарегистрированному потребителю собственную пропускную способность чтения 2 МБ/с независимо от других потребителей. Это устраняет ограничение чтения, когда несколько приложений используют один поток. Отслеживайте показатель GetRecords.IteratorAgeMilliseconds, чтобы обнаруживать отставание потребителей.
# Split shard to increase throughput
aws kinesis split-shard \
--stream-name clickstream \
--shard-to-split shardId-000000000001 \
--new-starting-hash-key 170141183460469231731687303715884105728
# Register an enhanced fan-out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/clickstream \
--consumer-name analytics-appKinesis Data Firehose: управляемая доставка
Kinesis Data Firehose — это полностью управляемая служба, которая собирает, преобразует и доставляет потоковые данные в такие места назначения, как S3, Amazon Redshift, Amazon OpenSearch Service, Splunk и конечные точки HTTP. Управлять сегментами не требуется — Firehose масштабируется автоматически. Настройте размер буфера (1–128 МБ) и интервал буфера (60–900 секунд); Firehose выполняет доставку, как только первым будет достигнут любой из этих лимитов.
# Create a Firehose delivery stream to S3
aws firehose create-delivery-stream \
--delivery-stream-name clickstream-to-s3 \
--s3-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"BucketARN": "arn:aws:s3:::my-data-lake-123",
"Prefix": "landing/clickstream/year=!{timestamp:yyyy}/month=!{timestamp:MM}/",
"BufferingHints": {"SizeInMBs": 64, "IntervalInSeconds": 300},
"CompressionFormat": "GZIP"
}'Преобразование данных Firehose с помощью Lambda
Firehose может вызывать функцию Lambda для каждой группы записей перед доставкой, чтобы преобразовывать, обогащать или фильтровать данные во время передачи. Типичные варианты использования: преобразование JSON в Parquet (с помощью схемы Glue), маскирование полей PII или удаление малозначимых событий. Записи, для которых преобразование завершилось ошибкой, можно отправлять в отдельный префикс ошибок S3 для повторной обработки, поэтому данные не теряются.
# Lambda transform function signature for Firehose
def lambda_handler(event, context):
output = []
for record in event['records']:
import base64, json
payload = json.loads(base64.b64decode(record['data']))
# Drop events with no user_id
if not payload.get('user_id'):
output.append({'recordId': record['recordId'], 'result': 'Dropped', 'data': record['data']})
else:
output.append({'recordId': record['recordId'], 'result': 'Ok', 'data': record['data']})
return {'records': output}Управляемая служба Apache Flink
Amazon Managed Service for Apache Flink (ранее Kinesis Data Analytics) запускает приложения Apache Flink в полностью управляемой инфраструктуре. Используйте её для аналитики с состоянием в реальном времени: агрегации в скользящих окнах, обнаружения аномалий, сопоставления шаблонов в последовательностях событий и объединения потоковых данных со справочными таблицами. Код Flink можно писать на Java, Python или Scala, а Flink управляет контрольными точками и состоянием с гарантией обработки ровно один раз.
# Flink SQL-style tumbling window (conceptual)
# Count page views per URL every 5 minutes
CREATE TABLE clickstream (
url STRING,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH ('connector' = 'kinesis', 'stream' = 'clickstream', ...);
SELECT
url,
COUNT(*) AS views,
TUMBLE_START(event_time, INTERVAL '5' MINUTE) AS window_start
FROM clickstream
GROUP BY url, TUMBLE(event_time, INTERVAL '5' MINUTE);Выбор между Streams и Firehose
Для экзамена SAA-C03 важно знать критерии выбора. Используйте Kinesis Data Streams, если необходимы задержка менее секунды, одновременное чтение несколькими потребителями или пользовательская логика обработки с полным контролем над хранением и повторным воспроизведением. Используйте Firehose, если нужно надёжно доставлять потоковые данные в S3, Redshift или OpenSearch с минимальным объёмом кода, дополнительным преобразованием во время передачи и автоматическим масштабированием, допускающим большую задержку (от 60 секунд).
Kinesis и SQS: классический экзаменационный выбор
В распространённом вопросе экзамена SAA-C03 предлагается выбрать между Kinesis и SQS. Основные различия: Kinesis сохраняет порядок сообщений внутри сегмента, поддерживает одновременное чтение одних и тех же данных несколькими потребителями и хранит записи для повторного воспроизведения. SQS удаляет сообщения после их получения (повторного воспроизведения нет), режим FIFO гарантирует строгий порядок, а сама служба лучше подходит для разделения микрослужб. Если в сценарии упоминаются аналитика в реальном времени или повторное воспроизведение, выбирайте Kinesis.
Производители Kinesis: SDK и KPL
Библиотека производителей Kinesis (KPL) — это высокопроизводительный клиент для записи данных в Kinesis Data Streams из приложений. KPL автоматически объединяет несколько небольших записей в один вызов API (до 1 МБ) и выполняет повторные попытки с увеличением интервала между ними. Это значительно снижает стоимость PUT для отдельных записей и увеличивает пропускную способность сегмента. Используйте KPL для производителей с большим объёмом данных, например веб-потоков кликов, датчиков IoT или конвейеров журналов.
# Basic KPL usage (Java pseudocode, no backticks)
KinesisProducer producer = new KinesisProducer();
byte[] data = 'hello world'.getBytes();
ByteBuffer buf = ByteBuffer.wrap(data);
// addUserRecord handles aggregation and retry internally
ListenableFuture future = producer.addUserRecord('clickstream', 'partitionKey', buf);
Futures.addCallback(future, new FutureCallback() { ... });Шаблон Firehose для Redshift
Распространённая архитектура — использовать Firehose как управляемый конвейер из потоков Kinesis или от непосредственных производителей в Amazon Redshift для создания хранилища данных. Сначала Firehose записывает данные во временное промежуточное хранилище S3, затем выполняет команду COPY, чтобы загрузить данные в Redshift. Это наиболее эффективный способ массовой загрузки потоковых данных в Redshift: прямые вставки строк по одной были бы крайне медленными из-за накладных расходов на каждую строку.
# Firehose Redshift destination (CLI snippet)
--redshift-destination-configuration '{
"RoleARN": "arn:aws:iam::123456789012:role/FirehoseRole",
"ClusterJDBCURL": "jdbc:redshift://cluster.xyz.us-east-1.redshift.amazonaws.com:5439/sales",
"CopyCommand": {
"DataTableName": "clickevents",
"CopyOptions": "JSON 'auto'"
},
"Username": "firehose_user",
"Password": "{{resolve:secretsmanager:redshift-pw}}",
"S3Configuration": {
"RoleARN": "...",
"BucketARN": "arn:aws:s3:::firehose-staging"
}
}'Мониторинг Kinesis с помощью CloudWatch
Основные показатели CloudWatch для Kinesis: IncomingBytes и IncomingRecords для измерения пропускной способности производителей, GetRecords.IteratorAgeMilliseconds для измерения отставания потребителей (высокое значение означает, что потребители не успевают обрабатывать данные) и WriteProvisionedThroughputExceeded для обнаружения превышения производителями лимитов сегментов. Настройте оповещения CloudWatch по возрасту итератора и превышению выделенной пропускной способности, чтобы запускать автоматическое масштабирование сегментов через Application Auto Scaling.
# CloudWatch alarm on high consumer lag
aws cloudwatch put-metric-alarm \
--alarm-name kinesis-high-lag \
--metric-name GetRecords.IteratorAgeMilliseconds \
--namespace AWS/Kinesis \
--dimensions Name=StreamName,Value=clickstream \
--statistic Maximum \
--period 60 \
--threshold 60000 \
--comparison-operator GreaterThanThreshold \
--evaluation-periods 3 \
--alarm-actions arn:aws:sns:us-east-1:123456789012:ops-alertsБыстрая проверка
Проверьте, насколько хорошо Вы усвоили концепции AWS Solutions Architect (SAA-C03) из этого урока.
Итоги урока
В этом уроке Вы узнали, что Kinesis Data Streams обеспечивает надёжную упорядоченную потоковую передачу с разделением на сегменты, несколькими потребителями и возможностью повторного воспроизведения, Kinesis Data Firehose предоставляет полностью управляемую доставку без кода в S3, Redshift и OpenSearch, а Managed Service for Apache Flink позволяет выполнять потоковую аналитику с состоянием в реальном времени. Далее мы рассмотрим стратегию миграции в облако «7 R».
Часто задаваемые вопросы
Урок «Потоки Kinesis, Firehose и аналитика в реальном времени» бесплатный?
Да — полный текст урока «Потоки Kinesis, Firehose и аналитика в реальном времени» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AWS Solutions Architect, подпишись на CoddyKit PRO. Курс AWS Solutions Architect содержит 4 уроков всего.
Чему я научусь в уроке «Потоки Kinesis, Firehose и аналитика в реальном времени»?
Получайте потоковые данные с помощью Kinesis Data Streams, передавайте их в S3 или Redshift через Firehose и анализируйте их в реальном времени с помощью Managed Service for Apache Flink. Ты практикуешь AWS Solutions Architect с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AWS Solutions Architect?
Предыдущий опыт не требуется. AWS Solutions Architect на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.
Сколько времени занимает урок «Потоки Kinesis, Firehose и аналитика в реальном времени»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AWS Solutions Architect?
Да. Каждый урок AWS Solutions Architect включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Создание озера данных на S3
- AWS Glue: ETL и каталог данных
- Amazon Athena: бессерверный SQL в S3
- Потоки Kinesis, Firehose и аналитика в реальном времени