流式数据
在 AWS 上开发应用
Ricardo Sueiras
Principal Technologist
流式数据基础
- 记录到达即持续处理。
- 生产者写入到流。
- 消费者独立并行读取。
- 吸收流量峰值,无需等待批处理。
流式数据由哪些部分组成?
- 生产者:发送记录的应用。
- 流:持久、有序的缓冲,解耦生产者与消费者。
- 消费者:读取并处理记录的应用。
- 存储与保留:定义记录可回放的保留时长。
分片(Shards)
- 分片是吞吐与并行处理的单位。
- 分区键将记录分布到各分片。
- 仅在分片内保证有序。
- 消费者用检查点跟踪进度并从故障恢复。
AWS 托管的流式数据服务
- Amazon Kinesis Data Streams。
- Amazon Data Firehose。
- 构建流式数据应用的核心基石。
Kinesis Data Streams
- 低延迟流式处理。
- 可配置保留期:24 小时(默认)至 365 天。
- 持久、可回放的流。
- 多个消费者可并行读取。
- 开发者用分片管理扩缩容。
何时使用 Kinesis Data Streams
- 需要实时低延迟时选择它。
- 需要记录回放。
- 有多个独立消费者。
- 需要自定义转换逻辑。
写入记录
- 分区键控制分片放置与顺序。
PutRecords 每次请求最多批量 500 条或 5 MB。
- 失败按记录报告。
FailedRecordCount 帮助定位失败记录。
消费记录:经典模式
- 消费者按需轮询 Kinesis API。
- 写入上限:1 MB/秒或 1,000 条/秒。
- 读取上限:2 MB/秒或每秒 5 次
GetRecords。
消费记录:起始位置
- 消费者指定从流的何处开始读取。
TRIM_HORIZON:分片中最早可用记录。
LATEST:下一个新记录。
AT_SEQUENCE_NUMBER 或 AFTER_SEQUENCE_NUMBER:从指定序列号。
AT_TIMESTAMP:在指定时间戳处或之后。
消费记录:增强扇出
- 每个已注册消费者在每个分片专享 2 MB/秒。
- 通过 HTTP/2 推送传递记录。
- 适用于多个独立消费者需低延迟时。
- 同一流还需高吞吐。
Lambda 事件源映射
- Lambda 轮询 Kinesis,并按批调用。
BatchSize 与 MaximumBatchingWindow 控制批次形态。
ParallelizationFactor:每个分片的并发调用数。
BisectBatchOnFunctionError:拆分失败批次以隔离问题记录。
ReportBatchItemFailures:避免重试已成功的记录。
OnFailure:将无法处理的批次发送至 SQS 或 SNS 的 DLQ。
Data Firehose
- 托管服务,将流式数据投递到目标。
- 目标包含 Amazon S3、Redshift 和 HTTP 端点。
- 支持缓冲、格式转换与 Lambda 转换。
- 延迟更高,缓冲数秒到数分钟。
- 无回放或长期保留。
- 适合快速将数据送达目标。
模式:热/冷
- 热/冷模式结合 Data Streams 与 Data Firehose。
- Data Streams 为生产者提供实时处理。
- Firehose 将同源记录归档到 S3。
Kinesis 与 SQS 对比
- 多消费者、需有序且可回放的流用 Kinesis。
- Kinesis 适合时间序列与有序数据。
- 任务分发用 SQS。
流式数据的扩展
- 数据流被拆分为分片。
- 分区键将记录路由到特定分片。
- 相同键始终落在同一分片,保持顺序。
- 分区键不均会导致热点分片。
- 用高基数字段缓解。
- 通过重分片扩缩容。
容量模式:预置
- 提供两种容量模式。
- 预置模式使用固定分片数,需手动管理。
- 成本与容量可预测。
- 需监控以避免限速与过度预置。
容量模式:按需
- 按需模式随流量自动扩缩。
- 无需手动管理分片。
- 适合负载难以预测的场景。
- 持续高流量下成本更高。
Kinesis Producer Library(KPL)
- AWS 库简化可扩展生产者与消费者的构建。
- Kinesis Producer Library(KPL)可批处理与聚合记录。
- 还可压缩记录并在失败时重试。
- 异步缓冲以最大化生产端吞吐。
Kinesis Consumer Library(KCL)
- Kinesis Consumer Library(KCL)负责分片协调。
- 管理工作器租约与检查点。
- 状态存于 DynamoDB 表。
- 消费者据此查找恢复位置。
故障处理:幂等性
- Kinesis 至少一次投递,可能出现重复。
- 将消费者设计为幂等。
- 用序列号或业务键在 DynamoDB 去重。
- 消费者可按序列号或时间戳重读保留记录。
故障处理:吞吐错误
ProvisionedThroughputExceededException 表示被限速。
- 通过增加分片数解决。
- 优化分区键分布。
- 带退避重试。
- 切换容量模式。
- 在 CloudWatch 跟踪
IncomingBytes 与 IncomingRecords。
- 关注
WriteProvisionedThroughputExceeded 以监控限速。
- 消费延迟发生于消费慢于生产时。
- 跟踪
GetRecords.IteratorAgeMilliseconds:值高表示滞后。
安全
- 使用 AWS KMS 的服务端静态加密。
- 传输中使用 TLS 加密。
- 通过 IAM 策略控制访问。
Preparing Video For Download...