0%
💡 왼쪽 원문을 읽으면서 오른쪽에 따라 써보세요. Tab 키로 힌트를 받을 수 있습니다.
1. Kinesis vs Apache Kafka 详细对比
1.1 架构差异
Apache Kafka 与 AWS Kinesis 都是面向实时数据流的平台, 但两者根本的架构理念并不相同。
Apache Kafka 架构
==================
+----------+ +------------------------------------------+
| Producer | --> | Kafka Cluster |
+----------+ | |
| Broker 1 Broker 2 Broker 3 |
| +-------+ +-------+ +-------+ |
| |Topic A| |Topic A| |Topic A| |
| |Part 0 | |Part 1 | |Part 2 | |
| | | | | | | |
| |Topic B| |Topic B| |Topic B| |
| |Part 0 | |Part 1 | |Part 2 | |
| +-------+ +-------+ +-------+ |
| |
| KRaft (Kafka 4.0+, 移除 ZooKeeper) |
+------------------------------------------+
|
v
+----------+
| Consumer |
| Group |
+----------+
AWS Kinesis Data Streams 架构
==============================
+----------+ +------------------------------------------+
| Producer | --> | Kinesis Stream (完全托管型) |
+----------+ | |
| Shard 1 Shard 2 Shard 3 |
| +-------+ +-------+ +-------+ |
| |Records| |Records| |Records| |
| |1MB/s W| |1MB/s W| |1MB/s W| |
| |2MB/s R| |2MB/s R| |2MB/s R| |
| +-------+ +-------+ +-------+ |
| |
| AWS 完全托管 (无需管理基础设施) |
+------------------------------------------+
|
v
+----------+
| Consumer |
| (KCL) |
+----------+
1.2 综合对比表
| 对比项目 | AWS Kinesis | Apache Kafka (自建运维) | Amazon MSK |
|---|---|---|---|
| 管理方式 | 完全托管型 | 自建运维 | 托管 Kafka |
| 扩展单位 | 分片 | 分区 + Broker | Broker |
| 每个分片/分区的写入 | 1 MB/s | 无限制 (取决于磁盘 I/O) | 依赖磁盘 I/O |
| 每个分片/分区的读取 | 2 MB/s (共享) | 无限制 | 依赖磁盘 I/O |
| 最长保留期 | 365 天 | 无限制 | 无限制 |
| 顺序保证 | 分片内 | 分区内 | 分区内 |
| 记录最大尺寸 | 1 MB | 默认 1 MB (可调整配置) | 默认 1 MB |
| 消费者模型 | KCL、增强型扇出 | 消费者组 | 消费者组 |
| 协议 | HTTPS/HTTP2 | 自有 TCP 协议 | 自有 TCP |
| 生态系统 | AWS 服务集成 | Kafka Connect、Schema Registry 等 | Kafka 生态系统 |
| 运维复杂度 | 低 | 高 | 中 |
| 初始搭建 | 数分钟 | 数小时 ~ 数天 | 数十分钟 |
1.3 吞吐量与延迟
吞吐量对比
===========
Kinesis (预置模式):
写入: 1 MB/s x 分片数
读取: 2 MB/s x 分片数 (共享)
2 MB/s x 分片数 x 消费者数 (增强型扇出)
示例: 100 个分片
写入: 100 MB/s
读取: 200 MB/s (共享) 或 200 MB/s x N (增强型扇出)
Kafka (自建运维):
取决于 Broker 性能
单个 Broker: 可达数百 MB/s
集群: 可处理数 GB/s
示例: 6 个 Broker 的集群
写入: 600+ MB/s
读取: 数 GB/s
延迟对比
=========
Kinesis:
- PutRecord: 数十 ms
- GetRecords (共享扇出): ~200 ms
- Enhanced Fan-Out: ~70 ms
Kafka:
- 生产者 -> 消费者: ~2-10 ms (依赖网络)
- 端到端: ~10-50 ms
1.4 成本对比
月度成本估算 (以美国东部为准)
================================
场景: 持续 10 MB/s 数据摄取,3 个消费者
Kinesis (预置模式):
- 需要 10 个分片 (10 MB/s / 1 MB/s per shard)
- 分片成本: 10 x 0.015 x 720 小时 = ~108 USD
- PUT 单元: ~360 USD (约 25.9B 单元/月)
- 增强型扇出 (3 个消费者): ~324 USD
合计: ~792 USD/月
Kinesis (按需 Advantage):
- 数据写入: ~25.9 TB x 0.032 = ~829 USD
- 数据读取: ~25.9 TB x 3 x 0.016 = ~1,243 USD
合计: ~2,072 USD/月
Kafka (EC2 自建运维):
- 3 个 Broker (m5.xlarge): 3 x 140 = ~420 USD
- EBS 存储 (1TB x 3): ~300 USD
- 运维人力成本: 另计
合计: ~720 USD/月 + 运维成本
Amazon MSK:
- 3 个 Broker (kafka.m5.large): ~456 USD
- 存储: ~300 USD
合计: ~756 USD/月
1.5 何时该选择什么
选择 Kinesis 的场景:
- 与 AWS 环境深度集成的架构
- 希望把运维负担降到最低时
- 需要与 Lambda、Firehose 等 AWS 服务直接联动时
- 中小规模吞吐量 (数十 MB/s 以下)
- 需要快速做原型验证时
选择 Kafka 的场景:
- 需要极高吞吐量时 (数 GB/s)
- 多云或混合云环境
- 需要 Kafka Connect 生态系统时
- 要求超低延迟时 (数 ms)
- 需要无限期保留数据时
2. Kinesis vs SQS:何时使用哪一个
| 对比项目 | Kinesis Data Streams | Amazon SQS |
|---|---|---|
| 数据处理模型 | 流式处理 (连续处理) | 消息队列 (逐条处理) |
| 消费者数量 | 多消费者同时处理 | 基本上为单消费者 |
| 顺序保证 | 分片内保证 | 仅 FIFO 队列保证 |
| 数据保留 | 24 小时 ~ 365 天 | 最长 14 天 |
| 数据重放 | 可以 (基于序列号) | 不可以 (处理后即删除) |
| 吞吐量 | 每个分片 1 MB/s 写入 | 几乎无限制 |
| 消息尺寸 | 最大 1 MB | 最大 256 KB |
| 延迟 | 毫秒级 | 毫秒级 |
| 计费模型 | 分片小时数 + 数据传输量 | 基于请求数 |
| 主要用途 | 实时分析、日志采集 | 微服务解耦 |
使用场景决策流程图
================================
数据处理需求分析
|
+----+----+
| |
同一份数据需要 消息只需要被
多个消费者读取 处理一次吗?
吗? |
| SQS
|
是否必须保证
实时顺序?
|
+----+----+
| |
YES NO
| |
Kinesis SQS FIFO
Data 或
Streams Kinesis
3. Amazon Managed Service for Apache Flink
3.1 概述
Amazon Managed Service for Apache Flink (原 Kinesis Data Analytics) 是可以在完全托管的 基础设施上运行 Apache Flink 的服务。
注意:原有的 Kinesis Data Analytics for SQL 自 2025 年 10 月起已停止新建, 官方推荐迁移到 Amazon Managed Service for Apache Flink。
3.2 主要功能
Managed Flink 架构
====================
+-----------+ +----------------------------+ +-----------+
| 源 | | Managed Flink | | 汇 |
| | --> | | --> | |
| - Kinesis | | +------------------------+ | | - Kinesis |
| - MSK | | | Flink Application | | | - S3 |
| - S3 | | | | | | - DynamoDB|
| | | | - SQL 查询 | | | - Open- |
| | | | - Java/Scala 应用 | | | Search |
| | | | - Python (PyFlink) | | | - Redshift|
| | | | | | | |
| | | | 窗口聚合 | | | |
| | | | 模式检测 | | | |
| | | | CEP (复合事件处理) | | | |
| | | +------------------------+ | | |
+-----------+ +----------------------------+ +-----------+
3.3 窗口处理类型
这是把流数据按时间分组后进行分析的核心功能。
窗口类型
=========
1) 滚动窗口 (Tumbling Window)
- 固定大小,互不重叠
|-------|-------|-------|-------|
0 5 10 15 20 (秒)
2) 滑动窗口 (Sliding/Hopping Window)
- 固定大小,按固定间隔滑动
|-----------|
|-----------|
|-----------|
0 2 4 6 8 10 (秒)
大小: 6 秒, 滑动: 2 秒
3) 会话窗口 (Session Window)
- 基于活动,以静默间隔划分
|---event-event---| gap |--event-event-event--| gap |
<-- Session 1 --> <------ Session 2 ---->
4) 全局窗口 (Global Window)
- 对整条流使用单一窗口
3.4 Flink SQL 示例
-- 定义 Kinesis 源表
CREATE TABLE clickstream (
user_id VARCHAR,
page VARCHAR,
action VARCHAR,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kinesis',
'stream' = 'clickstream-data',
'aws.region' = 'ap-northeast-2',
'scan.stream.initpos' = 'LATEST',
'format' = 'json'
);
-- 用 1 分钟滚动窗口按页面聚合浏览量
SELECT
page,
COUNT(*) AS view_count,
COUNT(DISTINCT user_id) AS unique_users,
TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(event_time, INTERVAL '1' MINUTE) AS window_end
FROM clickstream
WHERE action = 'view'
GROUP BY
page,
TUMBLE(event_time, INTERVAL '1' MINUTE);
-- 异常行为检测: 1 分钟内同一用户点击 10 次以上
SELECT
user_id,
COUNT(*) AS click_count,
TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start
FROM clickstream
WHERE action = 'click'
GROUP BY
user_id,
TUMBLE(event_time, INTERVAL '1' MINUTE)
HAVING COUNT(*) >= 10;
4. 实战流式架构模式
4.1 日志聚合流水线
日志聚合架构
==============
+----------+ +----------+ +---------+ +----------+
| App | | Kinesis | | Firehose| | S3 |
| Server 1 | --> | Agent | --> | | --> | (Raw |
+----------+ +----------+ | | | Logs) |
| | +----------+
+----------+ +----------+ | | |
| App | | Kinesis | | | v
| Server 2 | --> | Agent | --> | | +----------+
+----------+ +----------+ | | | Athena |
| | | (Query) |
+----------+ +----------+ | | +----------+
| App | | Kinesis | | |
| Server N | --> | Agent | --> | |
+----------+ +----------+ +---------+
|
v
+----------+ +----------+
| Lambda | --> | Open- |
| (实时 | | Search |
| 告警) | | (搜索/ |
+----------+ | 仪表板) |
+----------+
4.2 实时分析仪表板
实时仪表板架构
================
+----------+ +----------+ +-----------+ +----------+
| 网页/移动| | API | | Kinesis | | Managed |
| 客户端 | --> | Gateway | --> | Data | --> | Flink |
+----------+ +----------+ | Streams | | (聚合/ |
+-----------+ | 分析) |
+----------+
|
+---------+----+----+---------+
| | | |
v v v v
+--------+ +--------+ +--------+ +--------+
|DynamoDB| |Timestream| | S3 | |CloudWatch|
|(实时 | |(时序 | |(长期 | |(指标 |
| 数据) | | 数据) | | 存储) | | 告警) |
+--------+ +--------+ +--------+ +--------+
| |
v v
+--------------------+
| 仪表板应用 |
| (React/Vue + |
| WebSocket) |
+--------------------+
4.3 IoT 数据采集
# IoT 设备模拟器
import boto3
import json
import time
import random
from datetime import datetime
kinesis = boto3.client('kinesis', region_name='ap-northeast-2')
STREAM_NAME = 'iot-sensor-data'
def simulate_sensor(device_id):
"""IoT 传感器数据模拟"""
return {
'device_id': device_id,
'temperature': round(random.uniform(15.0, 45.0), 2),
'humidity': round(random.uniform(20.0, 90.0), 2),
'pressure': round(random.uniform(990.0, 1030.0), 2),
'battery_level': round(random.uniform(0.0, 100.0), 1),
'location': {
'lat': round(random.uniform(33.0, 38.0), 6),
'lon': round(random.uniform(126.0, 130.0), 6)
},
'timestamp': datetime.utcnow().isoformat() + 'Z'
}
def ingest_iot_data(num_devices=100, interval=1.0):
"""把 IoT 数据采集到 Kinesis"""
device_ids = [f'sensor-{i:04d}' for i in range(num_devices)]
while True:
records = []
for device_id in device_ids:
sensor_data = simulate_sensor(device_id)
records.append({
'Data': json.dumps(sensor_data).encode('utf-8'),
'PartitionKey': device_id
})
# PutRecords 最多 500 条
for batch_start in range(0, len(records), 500):
batch = records[batch_start:batch_start + 500]
response = kinesis.put_records(
StreamName=STREAM_NAME,
Records=batch
)
failed = response['FailedRecordCount']
if failed > 0:
print(f"Batch failed: {failed} records")
# 指数退避重试逻辑
retry_records = []
for i, result in enumerate(response['Records']):
if 'ErrorCode' in result:
retry_records.append(batch[i])
if retry_records:
time.sleep(0.5)
kinesis.put_records(
StreamName=STREAM_NAME,
Records=retry_records
)
print(f"Ingested {len(records)} sensor readings")
time.sleep(interval)
if __name__ == '__main__':
ingest_iot_data()
4.4 事件溯源模式
事件溯源架构
==============
+----------+ +----------+ +-----------+
| Command | | Kinesis | | Event |
| Handler | --> | Data | --> | Processor |
| | | Streams | | (Lambda/ |
| - 订单 | | (事件 | | ECS) |
| - 支付 | | 存储) | +-----------+
| - 配送 | +----------+ |
+----------+ | +-----+-----+
| | |
v v v
+----------+ +--------+ +--------+
| S3 | |DynamoDB| |SNS |
| (事件 | |(读取 | |(通知) |
| 归档) | | 模型) | | |
+----------+ +--------+ +--------+
事件流示例:
1. OrderCreated -> 订单创建事件
2. PaymentProcessed -> 支付处理事件
3. InventoryReserved -> 库存预留事件
4. ShipmentCreated -> 配送创建事件
4.5 ML 特征流水线
ML 特征流水线
===============
+----------+ +----------+ +-----------+ +----------+
| 事件 | | Kinesis | | Managed | | Feature |
| 源 | --> | Data | --> | Flink | --> | Store |
| | | Streams | | (特征 | | (Sage- |
| - 点击 | | | | 计算) | | Maker) |
| - 购买 | | | | | +----------+
| - 搜索 | | | | 实时: | |
+----------+ +----------+ | - 会话数 | v
| - 近期 | +----------+
| 购买数 | | ML 模型 |
| - 平均 | | 推理 |
| 停留时长 | +----------+
+-----------+
5. 性能优化
5.1 分区键设计
分区键设计是决定 Kinesis 性能最重要的因素。
好的分区键需要满足:
- 高基数 (唯一值必须足够多)
- 均匀分布 (不能倾斜到某个特定键上)
- 唯一键数量至少是分片数的 10 倍以上
好的分区键示例
================
1) 用户 ID (高基数)
user-001 -> Shard 1
user-002 -> Shard 3
user-003 -> Shard 2
...
分布均匀
2) UUID (分散度最佳)
随机 UUID -> 完美分散
缺点: 无法保证同一实体的顺序
3) 复合键
"region-userType-userId"
可以做精细的分散控制
不好的分区键示例
================
1) 日期 ("2026-03-20")
所有记录都落到同一个分片 -> 热点分片
2) 国家代码 ("KR", "US", "JP")
基数太低
流量会倾斜到特定国家
3) 固定值 ("default")
所有负载集中在单个分片
5.2 利用 KPL 进行聚合
KPL 聚合优化
=============
不使用聚合:
Record 1 (100B) -> PutRecord -> 1 次 API 调用
Record 2 (200B) -> PutRecord -> 1 次 API 调用
Record 3 (150B) -> PutRecord -> 1 次 API 调用
合计: 3 次 API 调用, 传输 450B
使用 KPL 聚合:
Record 1 (100B) --+
Record 2 (200B) --+--> 聚合记录 (450B) -> 1 次 API 调用
Record 3 (150B) --+
合计: 1 次 API 调用, 传输 450B
效果:
- API 调用次数大幅减少
- 降低 PUT 单元费用
- 吞吐量大幅提升
5.3 增强型扇出策略
# 注册增强型扇出消费者
import boto3
kinesis = boto3.client('kinesis', region_name='ap-northeast-2')
# 注册消费者
response = kinesis.register_stream_consumer(
StreamARN='arn:aws:kinesis:ap-northeast-2:123456789012:stream/my-stream',
ConsumerName='analytics-consumer'
)
consumer_arn = response['Consumer']['ConsumerARN']
print(f"Consumer ARN: {consumer_arn}")
# 确认消费者状态
response = kinesis.describe_stream_consumer(
StreamARN='arn:aws:kinesis:ap-northeast-2:123456789012:stream/my-stream',
ConsumerName='analytics-consumer'
)
print(f"Status: {response['ConsumerDescription']['ConsumerStatus']}")
5.4 错误处理与重试策略
import time
import random
def put_records_with_retry(kinesis_client, stream_name, records, max_retries=3):
"""使用指数退避的 PutRecords 重试"""
for attempt in range(max_retries):
response = kinesis_client.put_records(
StreamName=stream_name,
Records=records
)
failed_count = response['FailedRecordCount']
if failed_count == 0:
return response
# 只提取失败的记录
retry_records = []
for i, result in enumerate(response['Records']):
if 'ErrorCode' in result:
error_code = result['ErrorCode']
if error_code == 'ProvisionedThroughputExceededException':
retry_records.append(records[i])
else:
print(f"Non-retryable error: {error_code}")
if not retry_records:
return response
records = retry_records
# 指数退避 + 抖动
backoff = min(2 ** attempt * 0.1, 5.0)
jitter = random.uniform(0, backoff * 0.5)
wait_time = backoff + jitter
print(f"Retry {attempt + 1}: {len(retry_records)} records, waiting {wait_time:.2f}s")
time.sleep(wait_time)
print(f"Failed after {max_retries} retries: {len(records)} records")
return None
6. 监控:CloudWatch 指标
6.1 核心监控指标
| 指标 | 说明 | 告警阈值 |
|---|---|---|
| IncomingBytes | 进入流的字节数 | 分片容量的 80% |
| IncomingRecords | 进入流的记录数 | 每个分片 800 rec/s |
| GetRecords.IteratorAgeMilliseconds | 消费者落后了多少 | 60,000 ms (1 分钟) |
| WriteProvisionedThroughputExceeded | 写入超限次数 | 超过 0 即告警 |
| ReadProvisionedThroughputExceeded | 读取超限次数 | 超过 0 即告警 |
| GetRecords.Latency | GetRecords 调用延迟 | 1,000 ms |
| PutRecord.Latency | PutRecord 调用延迟 | 1,000 ms |
| GetRecords.Success | GetRecords 成功率 | 低于 99% 即告警 |
6.2 增强型监控
启用增强型监控后新增的指标
============================
分片级别指标:
- IncomingBytes (按分片)
- IncomingRecords (按分片)
- IteratorAgeMilliseconds (按分片)
- OutgoingBytes (按分片)
- OutgoingRecords (按分片)
- ReadProvisionedThroughputExceeded (按分片)
- WriteProvisionedThroughputExceeded (按分片)
热点分片探测:
Shard 1: IncomingBytes = 200 KB/s [正常]
Shard 2: IncomingBytes = 950 KB/s [警告! 已接近上限]
Shard 3: IncomingBytes = 300 KB/s [正常]
-> 建议拆分 Shard 2
7. 最佳实践与反模式
7.1 最佳实践
1) 分区键设计
- 使用高基数的键 (用户 ID、设备 ID)
- 确保唯一键数量在分片数的 10 倍以上
- 需要顺序时,把实体 ID 作为分区键
2) 生产者优化
- 使用 PutRecords (批量) API 把 API 调用次数降到最低
- 使用 KPL 优化记录聚合与采集
- 实现合适的重试逻辑 (指数退避 + 抖动)
3) 消费者优化
- 多消费者场景使用增强型扇出
- 使用 KCL 让分布式处理自动化
- 优化检查点周期 (过于频繁会增加 DynamoDB 费用)
4) 容量管理
- 可预测的负载:预置模式
- 不规则的负载:按需模式
- 用 CloudWatch 告警触发自动扩缩容
5) 成本优化
- 评估按需 Advantage 模式 (2025 年推出)
- 保留期只设置为实际需要的时长
- 注销不再需要的增强型扇出消费者
7.2 反模式
1) 使用单一分区键
- 所有数据集中到一个分片
- 立刻触达分片上限
2) 分片数量过多
- 成本增加,管理复杂度上升
- KCL 的 DynamoDB 租约表负载增加
3) 不使用检查点
- 故障时出现数据重复处理或丢失
- 务必实现恰当的检查点策略
4) 缺少错误处理
- 忽略 ProvisionedThroughputExceededException
- 不做重试导致失败数据丢失
5) 过度调用 GetRecords
- 必须遵守每个分片每秒 5 次的限制
- 设置合适的轮询间隔
8. 综合对比汇总表
| 项目 | Kinesis Data Streams | Kinesis Firehose | Kafka | SQS | Managed Flink |
|---|---|---|---|---|---|
| 类型 | 数据流式处理 | 数据投递 | 数据流式处理 | 消息队列 | 流式处理 |
| 管理 | 完全托管型 | 完全托管型 | 自建/MSK | 完全托管型 | 完全托管型 |
| 延迟 | ms | 60s+ | ms | ms | ms |
| 顺序保证 | 分片内 | 无 | 分区内 | 仅 FIFO | 取决于输入 |
| 重放 | 可以 | 不可以 | 可以 | 不可以 | 取决于输入 |
| 扩缩容 | 增加分片 | 自动 | 分区/Broker | 自动 | 增加 KPU |
| 转换 | 无 (由消费者负责) | Lambda | Kafka Streams | 无 | Flink 应用 |
| 计费模型 | 分片+数据 | 数据量 | 实例 | 请求数 | KPU 小时 |
| AWS 集成 | 高 | 非常高 | 低/中 | 高 | 高 |
| 最佳用途 | 实时采集 | 自动投递 | 大批量流式处理 | 任务队列 | 实时分析 |
9. 总结
服务选型指南
设计实时流式架构时,选择契合需求的服务至关重要。
- 如果目标只是简单的数据投递:使用 Amazon Data Firehose 直接投递到 S3、Redshift 等
- 如果需要实时处理与多消费者:Kinesis Data Streams + KCL 或增强型扇出
- 如果需要复杂的流式分析:Managed Service for Apache Flink
- 如果需要大容量处理与生态系统:Apache Kafka 或 Amazon MSK
- 如果需要微服务之间的消息通信:Amazon SQS
理解每个服务的强项,并且经常把多个服务组合起来使用,才是实战中的最优模式。 例如,用 Kinesis Data Streams 采集、用 Managed Flink 做实时分析、用 Data Firehose 长期存储到 S3,这样的组合是非常常见的架构。
현재 단락 (1/529)
Apache Kafka 与 AWS Kinesis 都是面向实时数据流的平台,
작성 글자: 0자원문 글자: 12,676자작성 단락: 0/529