0Pricing
Cloud & IT Cert Prep · 课时

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-app

Kinesis 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 在 S3 上构建数据湖
  2. AWS Glue:ETL 与数据目录
  3. Amazon Athena:S3 上的无服务器 SQL
  4. Kinesis Streams、Firehose 与实时分析
← 返回 Cloud & IT Cert Prep