流式传输执行模型
Amazon Glue 流式传输提供了两种用于处理流数据的执行模型。请选择最符合您的延迟要求和工作负载特性的模型。
微批处理模式(默认)
微批处理模式是所有 Amazon Glue 流式传输作业的默认执行模型。此模式使用 forEachBatch 或 Trigger.ProcessingTime 按配置的时间间隔轮询源。
在每个间隔期间,Amazon Glue 执行以下步骤:
规划执行 DAG。
启动任务。
从源中读取累积数据。
处理数据。
提交结果。
终止任务并重复。
由于存在每批次的调度开销,微批处理模式下的最低延迟通常为 1 至 2 秒。此模式支持所有来源(Kafka、Kinesis)、所有语言(Python、Scala)、有状态和无状态操作以及自动扩缩。
微批处理模式最适合大多数可接受秒级延迟的流式传输工作负载。
实时模式(Amazon Glue 6.0+)
实时模式是 Spark Structured Streaming 的一种新执行模型,从 Amazon Glue 6.0 开始提供,可将端到端延迟降低至亚秒级。实时模式还可帮助符合条件的工作负载实现毫秒级延迟。任务持续运行,在记录到达时对其进行处理,而不是等待数据累积。实时模式仅适用于 Spark Structured Streaming,不适用于旧版 Spark Streaming(DStreams)。
实时模式需要通过 --enable-real-time-mode 作业参数明确选择加入。此模式不使用 forEachBatch。相反,您可以直接将 writeStream 与 Trigger.RealTime 配合使用。
实时模式具有以下要求和限制:
来源:仅限 Kafka
操作:仅限无状态
语言:仅限 Scala
自动扩缩:不支持。请勿为实时模式作业启用自动扩缩。请使用固定的工作线程计数。
实时模式最适合低延迟的无状态转换,例如需要亚秒级延迟的 Kafka 到 Kafka 管道。
有关启用和使用实时模式的完整详细信息,请参阅 为流式传输作业启用实时模式。
执行模型比较
下表比较了这两种执行模型。
| 功能 | 微批处理模式 | 实时模式 |
|---|---|---|
| 延迟 | 秒到分钟 | 亚秒级 |
| 来源 | Kafka、Kinesis | 仅限 Kafka |
| 语言 | Python、Scala | 仅限 Scala |
| 操作 | 有状态和无状态 | 仅无状态 |
| 输出模式 | 追加、更新、完成 | 仅限更新 |
| 自动扩缩 | 是 | 否 |
| 触发器 | Trigger.ProcessingTime / forEachBatch | Trigger.RealTime |
| Amazon Glue 版本 | 所有 版本 | 6.0+ |