Kinesis Data Streams для обработки событий в реальном времени
Создавайте и обрабатывайте высокопроизводительные потоки событий с помощью Kinesis Data Streams, управляйте сегментами для обеспечения пропускной способности и используйте Lambda в качестве потребителя.
«Kinesis Data Streams для обработки событий в реальном времени» — бесплатный урок AWS Solutions Architect на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения AWS Solutions Architect, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс AWS Solutions Architect содержит 4 уроков всего.
Основные концепции Kinesis Data Streams
Kinesis Data Streams (KDS) — это долговечный потоковый сервис для работы с данными в реальном времени, сохраняющий порядок записей. Data организуются в поток, состоящий из одного или нескольких шардов. Каждый шард представляет собой упорядоченную последовательность записей данных. Производители помещают записи в шарды, используя ключ секционирования, который определяет, какой шард получит запись. Потребители считывают записи из шардов и обрабатывают их в порядке поступления в пределах каждого шарда.
Ограничения ёмкости и пропускной способности шардов
Каждый шард поддерживает пропускную способность записи 1 MB/s или 1 000 записей/s и пропускную способность чтения 2 MB/s (общую для всех стандартных потребителей этого шарда). Общая ёмкость потока линейно увеличивается с количеством шардов. Используйте формулу: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). Если производители достигают ограничения записи, появляются ошибки ProvisionedThroughputExceededException — устраните проблему, разделив шарды или распределив ключи секционирования более равномерно.
# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed
aws kinesis create-stream \
--stream-name iot-telemetry \
--shard-count 5Ключи секционирования и распределение данных
Ключ секционирования — это строка, которую Kinesis хеширует (MD5), чтобы определить, какой шард получит запись. Удачно выбранный ключ секционирования равномерно распределяет записи по шардам (предотвращение перегрузки отдельных шардов). Для IoT используйте идентификатор устройства. Для потоков кликов используйте идентификатор сеанса или пользователя. Избегайте ключей с малым количеством возможных значений (например, названия страны, если их всего 5), поскольку они приводят к перегрузке отдельных шардов: один шард получает непропорционально большой поток записей, а остальные простаивают.
import boto3, json
client = boto3.client('kinesis', region_name='us-east-1')
# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
StreamName='iot-telemetry',
Data=json.dumps(event),
PartitionKey='sensor-42' # high-cardinality -> even distribution
)Стандартные потребители и Enhanced Fan-Out
Стандартные потребители совместно используют пропускную способность чтения 2 MB/s на шард, выполняя опрос с помощью GetRecords. Если на одном шарде находятся 3 потребителя и каждому требуется 2 MB/s, их пропускная способность будет ограничена: всего им доступно 2 MB/s. Enhanced Fan-Out (EFO) предоставляет каждому зарегистрированному потребителю отдельный канал чтения со скоростью 2 MB/s через постоянное соединение HTTP/2 с передачей данных от сервера (SubscribeToShard). EFO увеличивает стоимость каждого часа работы потребителя с шардом, но полностью устраняет конкуренцию за чтение.
# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--consumer-name real-time-analytics
# List registered consumers
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetryLambda как потребитель Kinesis
Lambda нативно интегрируется с Kinesis Data Streams через сопоставление источника событий. Lambda опрашивает поток, считывает пакеты записей и вызывает Вашу функцию. Настройте BatchSize (от 1 до 10 000 записей), StartingPosition (TRIM_HORIZON — для самых старых, LATEST — для самых новых) и BisectBatchOnFunctionError, чтобы разделять пакеты, обработка которых завершилась ошибкой. Коэффициент параллелизации (от 1 до 10) позволяет Lambda запускать несколько параллельных вызовов на каждый шард, чтобы успевать обрабатывать быстро растущие потоки.
# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
--function-name ProcessIoTEvents \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--batch-size 100 \
--starting-position LATEST \
--parallelization-factor 5 \
--bisect-batch-on-function-error true \
--destination-config '{
"OnFailure": {
"Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
}
}'Библиотека клиентов Kinesis (KCL)
Библиотека клиентов Kinesis (KCL) — это прикладная платформа для создания надёжных потребителей Kinesis на Java (или на нескольких языках через MultiLangDaemon). KCL выполняет перечисление шардов, управляет арендой (распределяет шарды между экземплярами Worker), сохраняет контрольные точки прогресса в DynamoDB и корректно обрабатывает разделение и объединение шардов. Каждый Worker KCL обрабатывает один или несколько шардов, а KCL автоматически перераспределяет шарды при добавлении Worker или их сбоях. Для потребительских приложений в рабочей среде KCL предпочтительнее прямого опроса через SDK.
# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
# 'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();Хранение данных и повторное воспроизведение
Kinesis Data Streams по умолчанию хранит записи в течение 24 часов (срок можно увеличить до 7 дней или до 365 дней с помощью долгосрочного хранения за дополнительную плату). В отличие от SQS, использованные записи не удаляются после чтения — они остаются доступными до истечения срока хранения. Это позволяет нескольким потребителям независимо читать одни и те же записи, а также выполнять повторное воспроизведение, сбросив контрольную точку потребителя на более ранний порядковый номер — что особенно полезно для исправления ошибок или заполнения данных для новых сервисов.
# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
--stream-name iot-telemetry \
--retention-period-hours 168
# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name iot-telemetry \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100Режим ёмкости On-Demand
Kinesis Data Streams поддерживает два режима ёмкости. Режим Provisioned: Вы вручную управляете количеством шардов и платите за каждый час работы шарда. Режим On-Demand: Kinesis автоматически масштабирует ёмкость шардов в зависимости от входной пропускной способности (по умолчанию до 200 MB/s для записи и 400 MB/s для чтения), а Вы платите за каждый GB записанных и извлечённых данных. Режим On-Demand идеально подходит для переменчивых или непредсказуемых моделей трафика, когда Вы не хотите управлять масштабированием шардов.
# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--stream-mode-details StreamMode=ON_DEMANDГарантии порядка в пределах шардов
Kinesis гарантирует порядок в пределах одного шарда — записи с одинаковым ключом секционирования всегда попадают в один и тот же шард и считываются в порядке записи. Однако порядок между шардами не гарантируется. Если приложению требуется глобальный порядок всех записей, используйте один шард (это ограничит пропускную способность до 1 MB/s) или измените архитектуру так, чтобы порядок требовался только в пределах группы ключей секционирования (например, порядок для каждого устройства). Это распространённое экзаменационное отличие от SQS FIFO, который обеспечивает строгие дедупликацию и упорядочивание.
Безопасность: шифрование и VPC
Kinesis Data Streams шифрует записи на стороне сервера при хранении с помощью AWS KMS (CMK или ключа, управляемого AWS), если включено шифрование на стороне сервера. Все данные при передаче шифруются с помощью TLS. Для приложений, работающих в VPC и не передающих данные потока через общедоступный интернет, используйте конечную точку интерфейса VPC (PrivateLink) для Kinesis, чтобы трафик полностью оставался внутри магистральной сети AWS — это важно для регулируемых рабочих нагрузок.
# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
--stream-name iot-telemetry \
--encryption-type KMS \
--key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123Kinesis Data Streams, SQS и Kafka: сравнение
На экзамене SAA-C03 сравнивайте Kinesis Data Streams с альтернативами. Kinesis и SQS: Kinesis сохраняет порядок в пределах шарда и поддерживает чтение одних и тех же данных несколькими потребителями; SQS удаляет сообщения после чтения. Kinesis и MSK (Kafka): MSK — это управляемый Apache Kafka. Используйте его, когда Вам нужна совместимость с протоколом Kafka, расширенные настройки топиков или перенос существующей Kafka из локальной инфраструктуры. Используйте Kinesis для потоковой обработки в экосистеме AWS с более тесной интеграцией с Lambda, Firehose и Flink. Выбирайте Kinesis, если в вопросе явно не упоминаются Kafka или требования совместимости с Kafka.
Быстрая проверка
Проверьте своё понимание концепций AWS и архитектора решений (SAA-C03) из этого урока.
Итоги урока
В этом уроке Вы узнали, что Kinesis Data Streams предоставляет сегментированные, упорядоченные и долговечные потоки с хранением по умолчанию в течение 24 часов и возможностью повторного воспроизведения, Enhanced Fan-Out выделяет каждому потребителю 2 MB/s на шард, устраняя конкуренцию за чтение, а режим On-Demand автоматически масштабирует шарды для непредсказуемого трафика. Далее мы рассмотрим шаблоны хореографии и оркестрации в событийной архитектуре.
Часто задаваемые вопросы
Урок «Kinesis Data Streams для обработки событий в реальном времени» бесплатный?
Да — полный текст урока «Kinesis Data Streams для обработки событий в реальном времени» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс AWS Solutions Architect, подпишись на CoddyKit PRO. Курс AWS Solutions Architect содержит 4 уроков всего.
Чему я научусь в уроке «Kinesis Data Streams для обработки событий в реальном времени»?
Создавайте и обрабатывайте высокопроизводительные потоки событий с помощью Kinesis Data Streams, управляйте сегментами для обеспечения пропускной способности и используйте Lambda в качестве потребите… Ты практикуешь AWS Solutions Architect с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать AWS Solutions Architect?
Предыдущий опыт не требуется. AWS Solutions Architect на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.
Сколько времени занимает урок «Kinesis Data Streams для обработки событий в реальном времени»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке AWS Solutions Architect?
Да. Каждый урок AWS Solutions Architect включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- EventBridge: шина событий и правила
- Step Functions: оркестрация бессерверных рабочих процессов
- Kinesis Data Streams для обработки событий в реальном времени
- Шаблоны хореографии и оркестрации