Skip to content
Published on

[AWS] Kinesis 完全指南:实时流式数据处理的一切

分享
Authors

1. 什么是流式数据

1.1 批处理 vs 流式处理

传统的数据处理采用批处理(batch)方式,也就是把数据积攒一段时间后一次性处理。 但在现代应用中,实时响应已经成为必需。

批处理 (Batch Processing)
============================
[数据收集] --> [存储库] --> [周期性处理] --> [结果]
     |                           |
     +---- 数分钟 ~ 数小时 ------+

流式处理 (Stream Processing)
==================================
[数据生成] --> [] --> [即时处理] --> [结果]
     |                      |
     +---- 毫秒 ~ 数秒 -----+

1.2 流式数据的应用场景

  • 实时日志分析:即时收集服务器日志以进行异常检测
  • IoT 传感器数据:来自数百万台设备的遥测数据
  • 点击流分析:实时追踪用户行为以实现个性化
  • 金融交易监控:为欺诈检测而做的实时交易分析
  • 社交媒体信息流:实时趋势分析与情感分析

2. AWS Kinesis 家族概览

AWS Kinesis 是用于采集、处理、分析实时流式数据的一整套完全托管服务。

+------------------------------------------------------------------+
|                    AWS Kinesis 家族                                |
+------------------------------------------------------------------+
|                                                                    |
|  +------------------+  +------------------+  +------------------+  |
|  | Kinesis Data     |  | Amazon Data      |  | Managed Service  |  |
|  | Streams          |  | Firehose         |  | for Apache Flink |  |
|  |                  |  |                  |  |                  |  |
|  | 实时数据         |  | 传输管道         |  | 流式分析         |  |
|  | 流式传输         |  | (S3, Redshift    |  | (SQL, Java,      |  |
|  |                  |  |  OpenSearch)  |  |  Python, Scala)  |  |
|  +------------------+  +------------------+  +------------------+  |
|                                                                    |
|  +------------------+                                              |
|  | Kinesis Video    |                                              |
|  | Streams          |                                              |
|  |                  |                                              |
|  | 视频流式传输     |                                              |
|  | 采集与分析       |                                              |
|  +------------------+                                              |
+------------------------------------------------------------------+
服务主要用途数据类型延迟
Data Streams实时数据采集/处理记录(字节)实时(毫秒)
Data Firehose数据传输管道记录(字节)准实时(60 秒~)
Managed Flink流式分析/转换流数据实时(毫秒)
Video Streams视频采集/播放媒体帧实时(1~10 秒)

3. Kinesis Data Streams 深入解析

3.1 核心架构

Kinesis Data Streams 是大规模实时数据流式服务的核心。

                        Kinesis Data Stream
+----------------------------------------------------------------------+
|                                                                      |
|  Producer        Shard 1          Shard 2          Shard 3           |
|  --------     +----------+     +----------+     +----------+         |
|  |App A | --> |Record 1  | --> |Record 4  | --> |Record 7  |         |
|  |App B | --> |Record 2  | --> |Record 5  | --> |Record 8  |         |
|  |App C | --> |Record 3  | --> |Record 6  | --> |Record 9  |         |
|  --------     +----------+     +----------+     +----------+         |
|                    |                |                |                |
|                    v                v                v                |
|               Consumer A       Consumer A       Consumer A           |
|               Consumer B       Consumer B       Consumer B           |
+----------------------------------------------------------------------+

3.2 核心组成要素

流 (Stream)

  • 分片的逻辑分组
  • 保证数据记录顺序的单位

分片 (Shard)

  • 流的基本吞吐量单位
  • 写入:每秒 1MB 或 1,000 条记录
  • 读取:每秒 2MB(共享扇出),每消费者每秒 2MB(增强扇出)

记录 (Record)

  • 数据的基本单位
  • 由分区键、序列号、数据 blob(最大 1MB)组成

分区键 (Partition Key)

  • 决定记录被分配到哪个分片
  • 使用 MD5 哈希映射到分片
  • 拥有相同分区键的记录会按顺序存入同一个分片

序列号 (Sequence Number)

  • 自动分配给每条记录的唯一标识符
  • 保证记录在分片内的顺序
分区键哈希过程
====================

Partition Key: "user-123"
         |
         v
    MD5("user-123") = 0x7A3B...
         |
         v
    Hash Range: 0 ~ 2^128 - 1
         |
         v
    Shard 1: [0 ~ 2^127]          <-- 映射到这个范围
    Shard 2: [2^127 ~ 2^128 - 1]

3.3 生产者 (Producers)

把数据发送到 Kinesis 流的方式有好几种。

1) PutRecord / PutRecords API

这是最基础的方式。

import boto3
import json

kinesis = boto3.client('kinesis', region_name='ap-northeast-2')

# 发送单条记录
response = kinesis.put_record(
    StreamName='my-data-stream',
    Data=json.dumps({
        'event_type': 'page_view',
        'user_id': 'user-123',
        'page': '/products/laptop',
        'timestamp': '2026-03-20T10:30:00Z'
    }).encode('utf-8'),
    PartitionKey='user-123'
)
print(f"Shard ID: {response['ShardId']}")
print(f"Sequence Number: {response['SequenceNumber']}")

# 发送多条记录 (批量)
records = []
for i in range(100):
    records.append({
        'Data': json.dumps({
            'event_type': 'click',
            'user_id': f'user-{i % 10}',
            'element': 'buy_button',
            'timestamp': '2026-03-20T10:30:00Z'
        }).encode('utf-8'),
        'PartitionKey': f'user-{i % 10}'
    })

response = kinesis.put_records(
    StreamName='my-data-stream',
    Records=records
)
print(f"Failed records: {response['FailedRecordCount']}")

2) KPL (Kinesis Producer Library)

这是面向高性能生产的库。

  • 记录聚合 (Aggregation):把多条小记录打包成一条 Kinesis 记录
  • 记录收集 (Collection):把多条 Kinesis 记录打包进一次 PutRecords 调用
  • 自动重试:自动重新发送失败的记录
  • CloudWatch 指标:自动发布性能指标
KPL 聚合过程
==============

User Record A (100 bytes) --+
User Record B (200 bytes) --+--> Kinesis Record (500 bytes)
User Record C (200 bytes) --+         |
                                      v
User Record D (300 bytes) --+    PutRecords API
User Record E (400 bytes) --+--> (单次调用)
User Record F (100 bytes) --+

3) 直接使用 AWS SDK

可以通过各种语言的 AWS SDK 直接调用 API。

4) Kinesis Agent

这是安装在服务器上、把日志文件自动发送到流的代理程序。

3.4 消费者 (Consumers)

1) 共享扇出 (Shared Fan-Out) - GetRecords

这是默认的消费方式,每个分片 2MB/秒的读取吞吐量由所有消费者共享。

import boto3
import json
import time

kinesis = boto3.client('kinesis', region_name='ap-northeast-2')

# 获取分片迭代器
response = kinesis.get_shard_iterator(
    StreamName='my-data-stream',
    ShardId='shardId-000000000000',
    ShardIteratorType='LATEST'  # TRIM_HORIZON, AT_SEQUENCE_NUMBER, AT_TIMESTAMP 等
)
shard_iterator = response['ShardIterator']

# 轮询记录
while True:
    response = kinesis.get_records(
        ShardIterator=shard_iterator,
        Limit=100
    )

    for record in response['Records']:
        data = json.loads(record['Data'].decode('utf-8'))
        print(f"Partition Key: {record['PartitionKey']}")
        print(f"Sequence Number: {record['SequenceNumber']}")
        print(f"Data: {data}")

    shard_iterator = response['NextShardIterator']

    # GetRecords 每个分片每秒最多调用 5 次
    time.sleep(0.2)

2) 增强扇出 (Enhanced Fan-Out) - SubscribeToShard

这是使用 HTTP/2 的推送式投递方式。

共享扇出 vs 增强扇出
===============================

共享扇出 (Shared Fan-Out):
+--------+     +--------+
| Shard  | --> | 2MB/s  | --> Consumer A (轮询)
|        |     | (共享) | --> Consumer B (轮询)
+--------+     +--------+     Consumer C (轮询)

3 个消费者分摊使用 2MB/s
每个消费者: ~0.67 MB/s

增强扇出 (Enhanced Fan-Out):
+--------+     +--------+
| Shard  | --> | 2MB/s  | --> Consumer A (HTTP/2 推送)
|        |     | 2MB/s  | --> Consumer B (HTTP/2 推送)
|        |     | 2MB/s  | --> Consumer C (HTTP/2 推送)
+--------+     +--------+

每个消费者分配到专用的 2MB/s
延迟: ~70ms (共享: ~200ms+)

3.5 KCL (Kinesis Client Library)

KCL 是简化分布式消费者应用开发的库。

核心功能:

  • 租约管理:使用 DynamoDB 表追踪每个分片的所有权
  • 检查点:把处理进度保存到 DynamoDB,从而支持故障恢复
  • 自动负载均衡:按照工作节点数量自动重新分配分片
  • 重分片处理:分片拆分/合并时自动应对
KCL 架构
=============

                DynamoDB (Lease Table)
                +---------------------+
                | Shard ID | Worker   |
                |----------|----------|
                | shard-0  | worker-1 |
                | shard-1  | worker-2 |
                | shard-2  | worker-1 |
                | shard-3  | worker-2 |
                +---------------------+
                      ^         ^
                      |         |
              +-------+         +-------+
              |                         |
        +-----------+           +-----------+
        | Worker 1  |           | Worker 2  |
        | (EC2/ECS) |           | (EC2/ECS) |
        |           |           |           |
        | shard-0   |           | shard-1   |
        | shard-2   |           | shard-3   |
        +-----------+           +-----------+

3.6 数据保留期

  • 默认值:24 小时
  • 最大值:365 天
  • 保留期越长,成本越高
  • 需要数据重放时可以延长保留期

3.7 分片的拆分与合并

分片拆分 (Split)
==================
Shard 1 [0 ~ 100]
         |
    split at 50
         |
    +----+----+
    |         |
Shard 3    Shard 4
[0 ~ 50]  [51 ~ 100]

分片合并 (Merge)
==================
Shard 3 [0 ~ 50]     --> Shard 5
Shard 4 [51 ~ 100]   --> [0 ~ 100]
  • 拆分:用于分散热点分片的负载
  • 合并:用于把流量减少的两个相邻分片合到一起
  • 按需模式下会自动处理

3.8 容量模式

预置模式 (Provisioned Mode)

  • 直接指定分片数量
  • 适合可预测的工作负载
  • 按分片小时计费

按需模式 (On-Demand Mode)

  • 自动调整分片数量
  • 适合不可预测的工作负载
  • 按数据量计费
  • 2025 年新增:On-Demand Advantage 模式(数据使用费比原有模式便宜 60%)

4. Amazon Data Firehose(原 Kinesis Data Firehose)

4.1 概述

Amazon Data Firehose 是把流式数据稳定地传输到数据湖、数据存储、分析服务的完全托管服务。

Amazon Data Firehose 架构
================================

+----------+     +-----------+     +------------+     +----------+
| 数据源   |     | Firehose  |     | 转换       |     | 目标     |
|          | --> | 传输      | --> | (可选)     | --> |          |
| - Direct |     ||     | - Lambda   |     | - S3     |
| - Kinesis|     |           |     | - 格式转换 |     | - Redshift|
|   Data   |     | 缓冲:     |     | - 压缩     |     | - Open-  |
|   Streams|     | - 大小    |     | - 加密     |     |   Search |
| - MSK    |     | - 时间    |     |            |     | - Splunk |
|          |     |           |     |            |     | - HTTP   |
+----------+     +-----------+     +------------+     +----------+
                                                           |
                                                      +----------+
                                                      | 备份 S3  |
                                                      | (原始)   |
                                                      +----------+

4.2 核心特点

无需管理分片

  • 与 Data Streams 不同,不需要直接管理分片
  • 自动扩容/缩容

缓冲配置

  • 基于大小:1MB ~ 128MB
  • 基于时间:60 秒 ~ 900 秒
  • 两者中任意一个达到条件就会传输

数据转换

  • 通过 Lambda 函数做自定义转换
  • 自动转换为 Apache Parquet、ORC 等列式格式
  • 支持 Gzip、Snappy、Zip 压缩
  • SSE-S3 或 SSE-KMS 加密

4.3 Firehose 传输目标

目标特点
Amazon S3最常见,可转换为 Parquet/ORC
Amazon Redshift经由 S3 后用 COPY 命令加载
Amazon OpenSearch日志分析、搜索索引
Splunk安全监控、运维日志
HTTP 端点自定义目标,如 Datadog 等
Snowflake云数据仓库
Apache Iceberg开放表格式

4.4 代码示例:Firehose 发送

import boto3
import json

firehose = boto3.client('firehose', region_name='ap-northeast-2')

# 发送单条记录
response = firehose.put_record(
    DeliveryStreamName='my-firehose-stream',
    Record={
        'Data': json.dumps({
            'event_type': 'purchase',
            'user_id': 'user-456',
            'product': 'laptop',
            'amount': 1299.99,
            'timestamp': '2026-03-20T10:30:00Z'
        }).encode('utf-8')
    }
)

# 批量发送
records = []
for i in range(50):
    records.append({
        'Data': json.dumps({
            'event_type': 'page_view',
            'user_id': f'user-{i}',
            'page': f'/product/{i}',
            'timestamp': '2026-03-20T10:30:00Z'
        }).encode('utf-8')
    })

response = firehose.put_record_batch(
    DeliveryStreamName='my-firehose-stream',
    Records=records
)
print(f"Failed records: {response['FailedPutCount']}")

5. Kinesis Video Streams

5.1 概述

Kinesis Video Streams 是把摄像头、RADAR、LIDAR、无人机等产生的视频与媒体数据 安全地采集并可播放的完全托管服务。

5.2 主要功能

Kinesis Video Streams 架构
==================================

+----------+     +------------------+     +------------------+
| 设备     |     | Kinesis Video    |     | 消费者           |
|          | --> | Streams          | --> |                  |
| - 摄像头 |     |                  |     | - HLS 播放       |
| - 无人机 |     | - 自动扩展       |     | - DASH 播放      |
| - LIDAR  |     | - 持久化存储     |     | - GetMedia API   |
| - 智能   |     | - 加密           |     | - Rekognition    |
|   手机   |     |                  |     | - SageMaker      |
+----------+     +------------------+     +------------------+
  • HLS (HTTP Live Streaming):可在 Web 浏览器与移动端进行实时及归档播放
  • WebRTC:超低延迟的双向媒体流式传输
  • Amazon Rekognition 联动:人脸识别、物体检测等计算机视觉分析
  • 存储层级:热存储(实时访问)与温存储(成本高效的保管)

6. 计费模型

6.1 Kinesis Data Streams 费用

预置模式:

项目费用(以美国东部为准)
分片小时~0.015 USD/分片/小时
PUT 负载单元 (25KB)~0.014 USD/百万单元
增强扇出数据检索~0.013 USD/GB
增强扇出消费者分片小时~0.015 USD/消费者/分片/小时
长期保留(超过 24 小时)~0.023 USD/分片/小时

按需模式:

项目费用(以美国东部为准)
流小时 (Standard)~0.04 USD/小时
数据写入 (Standard)~0.08 USD/GB
数据读取 (Standard)~0.04 USD/GB
数据写入 (Advantage)~0.032 USD/GB
数据读取 (Advantage)~0.016 USD/GB

6.2 Amazon Data Firehose 费用

项目费用(以美国东部为准)
数据采集(前 500TB/月)~0.029 USD/GB
格式转换~0.018 USD/GB
VPC 传输~0.01 USD/GB + 按小时计费

7. 整体架构示例:实时分析管道

实时分析管道架构
=====================================

+----------+   +---------+   +----------+   +---------+   +----------+
| Web/     |   | Kinesis |   | Lambda   |   | Firehose|   | S3       |
| Mobile   |-->| Data    |-->| (转换/   |-->| 传输    |-->| (Data    |
| App      |   | Streams |   |  过滤)   |   ||   |  Lake)   |
+----------+   +---------+   +----------+   +---------+   +----------+
                    |                                           |
                    v                                           v
              +-----------+                              +----------+
              | Managed   |                              | Athena   |
              | Flink     |                              | (Ad-hoc  |
              | (实时     |                              |  查询)   |
              |  分析)    |                              +----------+
              +-----------+                                    |
                    |                                          v
                    v                                    +----------+
              +-----------+                              | Quick-   |
              | DynamoDB  |                              | Sight    |
              | (实时     |                              | (BI      |
              |  仪表板)  |                              |  仪表板) |
              +-----------+                              +----------+

7.1 完整的生产者/消费者示例

# === producer.py ===
import boto3
import json
import time
import random
from datetime import datetime

kinesis = boto3.client('kinesis', region_name='ap-northeast-2')
STREAM_NAME = 'clickstream-data'

def generate_click_event():
    """生成点击流事件"""
    pages = ['/home', '/products', '/cart', '/checkout', '/profile']
    actions = ['view', 'click', 'scroll', 'submit']
    user_id = f'user-{random.randint(1, 1000)}'

    return {
        'user_id': user_id,
        'page': random.choice(pages),
        'action': random.choice(actions),
        'session_id': f'sess-{random.randint(1, 100)}',
        'timestamp': datetime.utcnow().isoformat() + 'Z',
        'device': random.choice(['mobile', 'desktop', 'tablet']),
        'country': random.choice(['KR', 'US', 'JP', 'DE'])
    }

def send_events(batch_size=50, interval=1.0):
    """把事件发送到 Kinesis 流"""
    while True:
        records = []
        for _ in range(batch_size):
            event = generate_click_event()
            records.append({
                'Data': json.dumps(event).encode('utf-8'),
                'PartitionKey': event['user_id']
            })

        try:
            response = kinesis.put_records(
                StreamName=STREAM_NAME,
                Records=records
            )

            failed = response['FailedRecordCount']
            if failed > 0:
                print(f"Warning: {failed} records failed")

                # 重试失败的记录
                for i, record_response in enumerate(response['Records']):
                    if 'ErrorCode' in record_response:
                        print(f"  Error: {record_response['ErrorCode']}")
                        # 实现重试逻辑
            else:
                print(f"Sent {batch_size} records successfully")

        except Exception as e:
            print(f"Error: {e}")

        time.sleep(interval)

if __name__ == '__main__':
    send_events()
# === consumer.py ===
import boto3
import json
import time

kinesis = boto3.client('kinesis', region_name='ap-northeast-2')
STREAM_NAME = 'clickstream-data'

def get_shard_ids():
    """返回流中所有分片的 ID"""
    response = kinesis.describe_stream(StreamName=STREAM_NAME)
    return [
        shard['ShardId']
        for shard in response['StreamDescription']['Shards']
    ]

def process_records(records):
    """记录处理逻辑"""
    page_views = {}
    for record in records:
        data = json.loads(record['Data'].decode('utf-8'))
        page = data.get('page', 'unknown')
        page_views[page] = page_views.get(page, 0) + 1

        # 异常行为检测示例
        if data.get('action') == 'submit' and data.get('page') == '/checkout':
            print(f"[ALERT] Checkout event: user={data['user_id']}")

    for page, count in page_views.items():
        print(f"  Page: {page}, Views: {count}")

def consume_stream():
    """从流中读取数据并处理"""
    shard_ids = get_shard_ids()
    print(f"Found {len(shard_ids)} shards")

    shard_iterators = {}
    for shard_id in shard_ids:
        response = kinesis.get_shard_iterator(
            StreamName=STREAM_NAME,
            ShardId=shard_id,
            ShardIteratorType='LATEST'
        )
        shard_iterators[shard_id] = response['ShardIterator']

    while True:
        for shard_id in shard_ids:
            try:
                response = kinesis.get_records(
                    ShardIterator=shard_iterators[shard_id],
                    Limit=100
                )

                if response['Records']:
                    print(f"\n--- Shard: {shard_id} ---")
                    print(f"Records received: {len(response['Records'])}")
                    process_records(response['Records'])

                shard_iterators[shard_id] = response['NextShardIterator']

            except kinesis.exceptions.ExpiredIteratorException:
                # 迭代器过期时重新生成
                response = kinesis.get_shard_iterator(
                    StreamName=STREAM_NAME,
                    ShardId=shard_id,
                    ShardIteratorType='LATEST'
                )
                shard_iterators[shard_id] = response['ShardIterator']

        time.sleep(1)

if __name__ == '__main__':
    consume_stream()

8. 主要限制与注意事项

项目限制
记录最大大小1 MB
每次 PutRecords 请求最大记录数500 条
每次 PutRecords 请求最大大小5 MB
每个分片的写入吞吐量1 MB/秒 或 1,000 条记录/秒
每个分片的读取吞吐量(共享)2 MB/秒,GetRecords 每秒 5 次
每个分片的读取吞吐量(增强扇出)每个消费者 2 MB/秒
最大注册消费者数(增强扇出)每个流 20 个
数据保留24 小时(默认)~ 365 天
每个流的最大分片数默认 500(可申请提升)

9. 总结

AWS Kinesis 家族为实时数据流式处理提供了全面的解决方案。

  • Kinesis Data Streams:实时数据采集与处理的核心。基于分片的架构使其可扩展, 并通过 KCL 与增强扇出支持多种消费者模式。
  • Amazon Data Firehose:把数据自动传输到 S3、Redshift、OpenSearch 等的管道。 无需管理分片就能轻松使用。
  • Managed Flink:基于 Apache Flink 的托管流式分析服务, 同时支持 SQL 与代码两种分析方式。
  • Video Streams:面向视频/媒体数据采集、存储、播放的专用服务。

下一篇文章将讲解 Kinesis 的实战架构模式,以及与 Apache Kafka、SQS 的对比。