AWS Glue 流式处理
作为 AWS Glue 的一个组件,AWS Glue 流式处理使您能够近乎实时地高效处理流数据,以便您执行数据摄取、处理和机器学习等关键任务。借助 Apache Spark Streaming 框架,AWS Glue 流式处理提供了一种无服务器服务,可以大规模处理流数据。AWS Glue 在 Apache Spark 的基础上进行了各种优化,例如,无服务器基础设施、自动扩缩、可视化作业开发、流作业即时笔记本以及其他性能改进。
流式处理用例
AWS Glue 流式处理的一些常见用例如下:
近乎实时的数据处理:AWS Glue 流式处理使组织能够近乎实时地处理流数据,以便其根据最新信息获得见解并及时做出决策。
欺诈检测:您可以利用 AWS Glue 流式处理对流数据进行实时分析,这对于检测信用卡欺诈、网络入侵或网上诈骗等欺诈活动非常有用。通过持续处理和分析传入数据,您可以快速识别可疑模式或异常情况。
社交媒体分析:AWS Glue 流式处理可以处理实时社交媒体数据,比如推文、帖子或评论,使组织能够实时监控趋势、情绪分析和管理品牌声誉。
物联网(IoT)分析:AWS Glue 流式处理适用于处理和分析物联网设备、传感器和联网机器生成的高速数据流。它允许实时监控、异常检测、预测性维护和其他物联网分析用例。
点击流分析:AWS Glue 流式处理可以处理和分析来自网站或移动应用程序的实时点击流数据。这使企业能够深入了解用户行为,个性化用户体验,根据实时点击流数据优化营销活动。
日志监控和分析:AWS Glue 流式处理可持续实时处理和分析来自服务器、应用程序或网络设备的日志数据。这有助于检测异常、排查问题、监控系统运行状况和性能。
推荐系统:AWS Glue 流式处理以实时处理用户活动数据,动态更新推荐模型。这允许根据用户行为和偏好进行个性化和实时推荐。
以下是可以应用 AWS Glue 流式处理的各种用例的一些例子。它与 AWS 生态系统和托管服务集成,使其成为在云中进行实时流处理和分析的一个方便的选择。
使用 AWS Glue 流式处理有哪些好处?
使用 AWS Glue 流式处理的好处如下:
无服务器:AWS Glue 流式处理无服务器,无需管理基础设施。这减少了运营开销,使用户可以专注于数据处理和分析任务,而不是基础设施管理。
自动扩缩:AWS Glue 流式处理提供自动扩缩功能,可根据工作负载动态调整处理能力。它会自动扩展或缩减以处理数据量的波动,从而确保最佳性能和资源利用率。
视觉开发:流式处理作业开发可能很复杂。AWS Glue流式处理通过提供可视化创作工具 AWS Glue Studio 来应对这一挑战。AWS GlueStudio 简化了创建流式处理工作流的过程,使开发人员能够直观地设计和管理流应用程序,从而缩短学习曲线并提高工作效率。
经济高效:作为一项无服务器服务,AWS Glue 流式处理无需预置和维护基础设施,因而提高了成本效益。用户根据流式处理作业执行期间消耗的资源付费,从而根据实际使用量进行成本优化和扩缩。
处理复杂的工作负载:AWS Glue 流式处理专为处理复杂的流工作负载而设计。它可以处理和分析大量实时数据,支持高级转换,并与其他 AWS 服务集成,从而实现复杂的流式处理数据管道和分析工作流。
无锁定:AWS Glue 流式处理提供了灵活性,可避免供应商锁定。用户可将 AWS Glue 流式处理作为广泛的 AWS 生态系统的一部分,并与其他 AWS 服务无缝集成。这样就可以与现有的数据来源、应用程序和服务轻松集成,而不必受制于特定的技术或平台。
何时使用 AWS Glue 流式处理?
关于流式处理用例,有很多选择。我们建议在以下场景中使用 AWS Glue 流式处理。
如果您已经在使用 AWS Glue 或 Spark 进行批处理,那么 AWS Glue 流式处理是您的理想选择。它可以无缝过渡到构建流式处理作业,而无需学习新的语言或框架。AWS Glue 流式处理利用现有的知识和基础设施,简化了任务开发过程,使您能够轻松地将数据处理能力扩展到实时流场景。
如果您需要统一的服务或产品来处理批处理、流和事件驱动型工作负载,那么 AWS Glue 流式处理解决方案就是您的理想之选。有了 AWS Glue 流式处理您可以将数据处理需求整合到一个框架中,从而消除管理多个系统的复杂性。这样就能高效开发和维护各种数据工作流,同时确保不同工作负载类型之间的一致性和兼容性。
AWS Glue 流式处理非常适合涉及超大流数据量和复杂转换的场景,比如流式处理或关系数据库之间的连接。它可以高效处理和分析大量数据流,使您能够轻松处理要求苛刻的工作负载。无论是高速数据摄取还是复杂的数据操作,AWS Glue 流式处理的可扩展性和高级处理能力都能确保最佳性能和准确结果。
如果您更喜欢采用可视化方法来构建流式处理作业,AWS Glue 还提供了 AWS Glue Studio,您可以用它来直观地设计和管理您的流应用程序,从而简化开发过程。这种直观的界面使开发人员能够使用可视化界面创建、配置和监控流式处理工作流,从而缩短学习曲线并提高工作效率。
对于 SLA(服务水平协议)要求严格、超过 10 秒的近实时用例,AWS Glue 流式处理是一个极佳的选择。
如果您使用 Apache Iceberg、Apache Hudi 或 Delta Lake 构建事务数据湖,AWS Glue 流式处理为这些开放表格式提供了本机支持。这种无缝集成使您能够直接处理来自这些事务数据湖的流式处理数据,从而确保数据一致性、完整性和兼容性。
当需要为各种数据目标摄取流数据时:AWS Glue 流式处理为各种数据目标提供了本机目标,例如 Amazon Redshift、Amazon RDS、Amazon Aurora、Oracle、SQL Server 和其他目标。
支持的数据来源
AWS Glue 流式处理支持以下数据来源:
Amazon Kinesis
Amazon MSK(Managed Streaming for Apache Kafka)
自行管理的 Apache Kafka
支持的数据目标
AWS Glue 流式处理支持多个数据目标:
AWS Glue Data Catalog 支持的数据目标
Amazon S3
Amazon Redshift
MySQL
PostgreSQL
Oracle
Microsoft SQL Server
Snowflake
任何可以使用 JDBC 连接的数据库
Apache Iceberg、Delta 和 Apache Hudi
AWS Glue Marketplace 连接器
为流式传输作业启用实时模式
实时模式(RTM)是适用于 Spark Structured Streaming 的新执行模型,已在 AWS Glue 6.0 中推出。RTM 可将端到端延迟从数秒或数分钟缩短至亚秒级。实时模式仅适用于 Spark Structured Streaming 作业。它不适用于旧版 Spark Streaming(DStreams)或其他作业类型。
RTM 使用 Trigger.RealTime。任务在批处理窗口(默认为 5 分钟)内持续运行,并在记录到达时立即进行处理,而非跨时间间隔累积数据。这与默认的微批处理模型不同,在该模型中,forEachBatch/Trigger.ProcessingTime 每隔一段时间就会轮询、处理、提交并重新启动任务。
重要
RTM 需要通过作业参数显式选择加入。如果没有足够的任务槽来覆盖所有源分区,RTM 会静默地丢弃未分配的分区。您必须预置足够的工作线程以覆盖所有 Kafka 分区。
先决条件
在启用实时模式之前,请确认您的作业满足以下要求:
-
AWS Glue 版本 6.0
-
作业必须使用 Spark Structured Streaming。实时模式不适用于旧版 Spark Streaming(DStreams)或其他作业类型。
-
作业类型必须是 Spark Streaming(
gluestreaming命令) -
作业语言必须是 Scala(
--job-language scala)。PySpark RTM 支持在 Spark 4.2 之前的版本中不可用。 -
仅限 Kafka 来源。AWS Glue 6.0 中的 RTM 不支持 Amazon Kinesis。
-
仅限无状态操作(选择、筛选、投影、映射)。不支持有状态操作,例如聚合、联接、重复数据删除和窗口化操作。
-
输出模式必须为“更新”。RTM 不支持追加模式。
-
自动扩缩与实时模式不兼容。请勿为 RTM 作业启用自动扩缩。配置固定数量的工作线程,足以覆盖源主题中的所有 Kafka 分区。
何时使用实时模式
实时模式专为特定类别的流式传输工作负载而设计。请在以下情况下考虑使用实时模式:
-
您需要亚秒级的端到端延迟,而微批处理延迟(1–2 秒或更长时间)对于您的应用场景而言过高。
-
您的管道执行无状态转换,例如筛选、投影、扩充记录,或将记录从 Kafka 路由到 Kafka 或其他接收器。
-
您拥有数量固定且可预测的 Kafka 分区,并且可以相应地预置工作线程。
-
您的作业使用 Scala 编写。
在以下情况下,请继续使用微批处理模式:
-
您需要有状态操作,例如聚合、联接、重复数据删除或窗口计算。
-
您使用 Amazon Kinesis 作为源。
-
您编写 PySpark 作业。
-
您依靠自动扩缩功能来处理可变的数据量。
-
您使用
forEachBatch或 GlueContext 流式传输 API。 -
秒级延迟对于您的应用场景而言是可接受的。
实时模式的工作原理
以下介绍微批处理模型与实时模式之间的区别:
- 微批处理模式
-
每个间隔都会启动任务、读取累积的数据、处理数据、提交检查点、终止任务,然后重复。最小延迟约为 1–2 秒。
- 实时模式
-
任务启动一次,运行持续时间为
batchDurationMs(默认为 5 分钟)。任务在记录到达时即对其进行处理,延迟低于一秒。到达截止时间时,各任务协同停止。驱动程序提交检查点,下一批次重新启动任务。
两种模式均使用相同的检查点格式和恢复机制。关键区别在于任务生命周期。微批处理模式每隔一段时间就会终止并重新启动任务。实时模式使任务在较长的批处理窗口内连续运行。
重要
如果没有足够的任务槽来处理所有源分区,RTM 会静默地丢弃未分配的分区。确保预置足够的工作线程以覆盖所有分区。
要启用实时模式
您可以通过将 --enable-real-time-mode 作业参数设置为 true 来启用实时模式。您可以在 AWS Glue 控制台中或通过 API 设置此参数。
要启用实时模式(控制台)
-
打开 AWS Glue 控制台
,然后打开您的流式传输作业。 -
选择 Job details(任务详细信息)选项卡。
-
对于 Glue 版本,选择 Glue 6.0。对于类型,选择 Spark Streaming。
-
滚动到作业参数部分。
-
选择添加新参数。
-
对于键,输入
--enable-real-time-mode。对于 Value(值),输入true。 -
选择保存。
注意
前导破折号为必填项。作业参数是 DefaultArguments 的控制台视图。
要启用实时模式(API)
--enable-real-time-mode 标志存储在作业定义的 DefaultArguments 映射中。您可以在创建或更新作业时设置该标志。
要创建新作业(AWS CLI)
运行以下命令:
aws glue create-job \ --name my-rtm-job \ --role arn:aws:iam::123456789012:role/MyGlueRole \ --glue-version 6.0 \ --worker-type G.1X --number-of-workers 4 \ --command '{"Name":"gluestreaming","ScriptLocation":"s3://my-bucket/scripts/rtm-job.scala"}' \ --default-arguments '{ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/" }' \ --region us-east-2
要创建新作业(boto3)
使用以下代码:
import boto3 glue = boto3.client("glue", region_name="us-east-2") glue.create_job( Name="my-rtm-job", Role="arn:aws:iam::123456789012:role/MyGlueRole", GlueVersion="6.0", WorkerType="G.1X", NumberOfWorkers=4, Command={ "Name": "gluestreaming", "ScriptLocation": "s3://my-bucket/scripts/rtm-job.scala", }, DefaultArguments={ "--enable-real-time-mode": "true", "--job-language": "scala", "--class": "GlueApp", "--TempDir": "s3://my-bucket/tmp/", }, )
要更新现有作业(AWS CLI)
运行以下命令:
aws glue update-job \ --job-name my-existing-job \ --job-update '{ "GlueVersion": "6.0", "DefaultArguments": { "--enable-real-time-mode": "true", "--job-language": "scala" } }'
编写流式传输脚本
作业参数声明了使用实时模式的意图。您的脚本会选择触发器。
以下 Scala 示例显示了使用 Trigger.RealTime 的流式传输查询:
import org.apache.spark.sql.streaming.Trigger val query = df.writeStream .format("kafka") .outputMode("update") .trigger(Trigger.RealTime(60000L)) // checkpoint interval in milliseconds .start() query.awaitTermination()
Trigger.RealTime 以毫秒为单位设置检查点间隔。更新输出模式为必填项。追加模式会抛出 OUTPUT_MODE_NOT_SUPPORTED。
只要设置了该标志,您就可以在一个脚本中混合使用多种模式:
dfA.writeStream.outputMode("update").trigger(Trigger.RealTime(60000L)).start() dfB.writeStream.outputMode("append").trigger(Trigger.ProcessingTime("30 seconds")).start()
缺少标志时的行为
以下内容描述了作业在未设置 --enable-real-time-mode 标志时的行为:
-
启动不带
--enable-real-time-mode标志的实时查询的作业会在查询启动时失败。失败消息会指示您添加该参数。 -
仅限微批处理的作业不会因为缺少此标志而受到影响。
-
设置了该标志但仅使用微批量查询的作业也不受影响。
注意事项和限制
在使用实时模式时,请注意以下几点:
- 分区删除
-
如果任务槽数量不足以覆盖所有源分区,则未分配的分区将不会被处理。预置工作线程以覆盖所有 Kafka 分区。
- 无自动扩缩
-
请勿为实时模式作业启用自动扩缩。自动扩缩与 RTM 不兼容,并且引入了延迟,从而抵消了低延迟的益处。预置固定数量的工作线程,其数量应等于或大于源主题中的 Kafka 分区数。
- 仅限 Kafka
-
Amazon Kinesis 数据来源在 AWS Glue 6.0 中不支持 RTM。
- 仅限 Scala
-
在 Spark 4.2 版本之前,RTM 不支持 PySpark。
- 仅无状态
-
不支持聚合、联接、重复数据删除、窗口化操作和
transformWithState。 - forEachBatch 不兼容
-
RTM 不使用
forEachBatch模型。直接将writeStream与Trigger.RealTime配合使用。 - 检查点恢复
-
作业重新启动时,RTM 会从最后一个检查点恢复。检查点每隔
batchDurationMs触发一次。最坏情况下的再处理是一个批处理窗口的持续时间(至少一次语义)。