Kinesis Streams、Firehose 与实时分析
使用 Kinesis Data Streams 摄取流式数据,通过 Firehose 将其传送到 S3 或 Redshift,并使用 Managed Service for Apache Flink 进行实时分析
Kinesis Streams、Firehose 与实时分析 是 CoddyKit 上的免费 Cloud & IT Cert Prep 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Cloud & IT Cert Prep 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Cloud & IT Cert Prep 课程共包含 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 Data Stream 是一种持久、有序,并按分片进行分区的日志。每个分片提供 1 MB/s 的写入吞吐量和 2 MB/s 的读取吞吐量。数据记录默认保留 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'分片、吞吐量和扩展
分片数量决定流的总吞吐量。您可以拆分一个分片以使吞吐量翻倍,也可以合并两个分片以降低成本。使用 Enhanced Fan-Out 可让每个已注册的消费者独立获得 2 MB/s 的读取吞吐量,不受其他消费者影响,从而在多个应用程序消费同一数据流时避免读取限流。监控 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 MB)和缓冲区间隔(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"
}'使用 Lambda 转换 Firehose 数据
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 应用程序。它适用于有状态的实时分析,例如滑动窗口聚合、异常检测、事件序列模式匹配,以及将流数据与参考表连接。您可以使用 Java、Python 或 Scala 编写 Flink 代码,而 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。当您只需以最少的代码将流数据可靠地传送到 S3、Redshift 或 OpenSearch,并需要可选的传输中转换和自动扩展,同时可以接受较高延迟(60 秒以上)时,请使用 Firehose。
Kinesis 与 SQS:经典考试选择题
常见的 SAA-C03 题目会要求您在 Kinesis 和 SQS 之间进行选择。主要区别包括:Kinesis 保证分片内的消息顺序,支持多个消费者同时读取相同数据,并保留记录以便重放。SQS 在消息被消费后会将其删除(无法重放),FIFO 模式可保证严格的顺序,并且更适合解耦微服务。如果场景提到实时分析或重放,请选择 Kinesis。
Kinesis 生产者:SDK 和 KPL
Kinesis Producer Library (KPL) 是一个高吞吐量客户端库,用于从应用程序向 Kinesis Data Streams 写入数据。KPL 会自动将多个小记录聚合为一次 API 调用(最多 1 MB),并通过退避机制处理重试。这会大幅降低每条记录的 PUT 成本,并提高每个分片的吞吐量。对于 Web 点击流、物联网传感器或日志管道等高容量生产者,请使用 KPL。
# 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 的最高效方式;由于每行插入都会产生额外开销,直接逐行插入 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"
}
}'使用 CloudWatch 监控 Kinesis
Kinesis 的关键 CloudWatch 指标包括:用于衡量生产者吞吐量的 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 Streams、Firehose 与实时分析」课时是免费的吗?
是的 — 「Kinesis Streams、Firehose 与实时分析」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Cloud & IT Cert Prep 课程的其余内容,请升级到 CoddyKit PRO。 Cloud & IT Cert Prep 课程共包含 4 节课。
「Kinesis Streams、Firehose 与实时分析」这节课中我会学到什么?
使用 Kinesis Data Streams 摄取流式数据,通过 Firehose 将其传送到 S3 或 Redshift,并使用 Managed Service for Apache Flink 进行实时分析 你通过在浏览器中直接运行的动手代码来练习 Cloud & IT Cert Prep,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Cloud & IT Cert Prep 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Cloud & IT Cert Prep 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「Kinesis Streams、Firehose 与实时分析」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Cloud & IT Cert Prep 课中编写并运行代码吗?
能。每节 Cloud & IT Cert Prep 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 在 S3 上构建数据湖
- AWS Glue:ETL 与数据目录
- Amazon Athena:S3 上的无服务器 SQL
- Kinesis Streams、Firehose 与实时分析