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 的对比。
현재 단락 (1/485)
传统的数据处理采用批处理(batch)方式,也就是把数据积攒一段时间后一次性处理。