用于实时事件处理的 Kinesis Data Streams
使用 Kinesis Data Streams 生产和消费高吞吐量事件流,管理分片以满足吞吐量需求,并使用 Lambda 作为消费者
用于实时事件处理的 Kinesis Data Streams 是 CoddyKit 上的免费 Cloud & IT Cert Prep 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Cloud & IT Cert Prep 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Cloud & IT Cert Prep 课程共包含 4 节课。
Kinesis Data 流核心概念
Kinesis Data 流 (KDS) 是一种持久、有序的实时数据流服务。数据会被组织到一个由一个或多个分片组成的流中。每个分片都是数据记录的有序序列。生产者使用分区键将记录写入分片,分区键决定记录接收于哪个分片。消费者从分片读取记录,并按照记录在各分片中的到达顺序进行处理。
分片容量和吞吐量限制
每个分片支持1 MB/s 或每秒 1,000 条记录的写入吞吐量,以及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),以确定记录接收于哪个分片。经过合理选择的分区键可以将记录均匀分布到各个分片中,从而防止分片过热。对于物联网数据,请使用设备 ID。对于点击流,请使用会话 ID 或用户 ID。请避免使用低基数键(例如只有 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
标准消费者通过轮询使用 GetRecords,共享每个分片 2 MB/s 的读取吞吐量。如果一个分片上有 3 个消费者,且每个消费者都需要 2 MB/s,那么它们共享总计 2 MB/s 的吞吐量,并会受到限流。Enhanced Fan-Out (EFO) 会通过持久的 HTTP/2 推送连接(SubscribeToShard)为每个已注册的消费者提供专用的 2 MB/s 读取通道。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 流。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 Client Library (KCL)
Kinesis Client Library (KCL) 是用于构建可靠 Java(或通过 MultiLangDaemon 支持多语言)Kinesis 消费者的应用框架。KCL 负责分片枚举、租约管理(在多个 Worker 实例之间分配分片)、将进度检查点保存到 DynamoDB,以及优雅地处理分片拆分和合并。每个 KCL Worker 处理一个或多个分片,并且当 Worker 扩展或发生故障时,KCL 会自动重新平衡分片。对于生产环境中的消费者应用,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 流默认将记录保存 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 100On-Demand 容量模式
Kinesis Data 流支持两种容量模式。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 流会使用 AWS KMS(CMK 或 AWS 托管密钥)对记录进行服务器端静态加密。所有传输中的数据都会使用 TLS 加密。对于在 VPC 中运行且不应通过公共互联网发送流数据的应用,请为 Kinesis 使用VPC 接口端点(PrivateLink),这样流量就能完全留在 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 流、SQS 与 Kafka 的比较
在 SAA-C03 考试中,请将 Kinesis Data 流与其他方案进行比较。Kinesis 与 SQS:Kinesis 保证分片内的顺序,并支持多个消费者读取相同的数据;SQS 会在消息被消费后删除消息。Kinesis 与 MSK(Kafka):MSK 是托管式 Apache Kafka;当您需要 Kafka 协议兼容性、高级主题配置,或正在从本地 Kafka 迁移时,应使用 MSK。对于与 Lambda、Firehose 和 Flink 紧密集成的 AWS 原生数据流处理,应使用 Kinesis。除非题目明确提到 Kafka 或 Kafka 兼容性要求,否则请选择 Kinesis。
快速检查
测试您对本课 AWS Solutions Architect (SAA-C03) 概念的理解。
课程回顾
在本课中,您学习了:Kinesis Data 流提供分片化、有序且持久的数据流,默认保留 24 小时并支持重放;Enhanced Fan-Out 为每个消费者提供每个分片专用的 2 MB/s 吞吐量,从而消除读取竞争;以及On-Demand 模式会自动扩展分片,以应对难以预测的流量。接下来,我们将学习事件驱动架构中的 Choreography 与 orchestration 模式。
常见问题解答
「用于实时事件处理的 Kinesis Data Streams」课时是免费的吗?
是的 — 「用于实时事件处理的 Kinesis Data Streams」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Cloud & IT Cert Prep 课程的其余内容,请升级到 CoddyKit PRO。 Cloud & IT Cert Prep 课程共包含 4 节课。
「用于实时事件处理的 Kinesis Data Streams」这节课中我会学到什么?
使用 Kinesis Data Streams 生产和消费高吞吐量事件流,管理分片以满足吞吐量需求,并使用 Lambda 作为消费者 你通过在浏览器中直接运行的动手代码来练习 Cloud & IT Cert Prep,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Cloud & IT Cert Prep 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Cloud & IT Cert Prep 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。
「用于实时事件处理的 Kinesis Data Streams」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Cloud & IT Cert Prep 课中编写并运行代码吗?
能。每节 Cloud & IT Cert Prep 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- EventBridge:事件总线与规则
- Step Functions:编排无服务器工作流
- 用于实时事件处理的 Kinesis Data Streams
- 编舞模式与编排模式