View a markdown version of this page

从 REST API 读取数据 - AWS Glue

从 REST API 读取数据

注册 REST API ConnectionType 并创建 AWS Glue 连接后,您可以在 AWS Glue ETL 作业中从 REST API 读取数据。通过此连接,您可以在同一作业中将外部 REST API 数据与其他数据来源一起处理。您需要连接名称和实体名称才能读取数据。

以下示例演示如何使用 Python 从 REST API 数据来源读取数据:

rest_read = glueContext.create_dynamic_frame.from_options( connection_type="rest", connection_options={ "connectionName": "connection-name", "ENTITY_NAME": "entity-name", "CONNECTION_TYPE": "REST-connection-type" } )

筛选数据

您可以将筛选条件谓词下推至源 REST API,以减少传输的数据量。REST API 连接器支持两种筛选模式:

  • QUERY_PARAMS:每个筛选条件均作为单独的 URL 查询参数。例如:?created[gte]=1704067200&created[lte]=1717200000

  • FILTER_STRING:所有筛选条件合并为一个查询参数。例如:?search=status eq "ACTIVE" and lastUpdated gt "2024-01-01"

要在 AWS Glue ETL 作业中应用筛选条件谓词,请使用 FILTER_PREDICATE 连接选项:

rest_read = glueContext.create_dynamic_frame.from_options( connection_type="rest", connection_options={ "connectionName": "connection-name", "ENTITY_NAME": "entity-name", "CONNECTION_TYPE": "REST-connection-type", "FILTER_PREDICATE": "status = \"ACTIVE\" AND lastUpdated >= \"2024-01-01T00:00:00.000Z\"" } )
FilterConfiguration 属性

使用 FilterConfiguration 对象在 ConnectionType 注册中配置筛选行为。下表描述了可用的属性:

属性

类型

说明

FilterMode

字符串

必填项。QUERY_PARAMSFILTER_STRING

OperatorMappings

Map

将逻辑运算符映射到特定于 API 的语法。

DateTimeFormat

字符串

DateTime 格式模式或 EPOCH_SECONDS/EPOCH_MILLIS

StripQuotes

布尔值

指定是否从输入值中去除周围的引号。默认值:true

BetweenConfiguration

对象

默认 BETWEEN 处理配置。

FilterStringConfiguration

对象

特定于 FILTER_STRING 模式的设置。

FilterStringConfiguration properties

下表描述了 FILTER_STRING 模式的属性:

属性

类型

说明

FilterStringKey

字符串

必需。查询参数键(例如,“search”或“filter”)。

QuoteStringValues

布尔值

指定是否将 String 和 DateTime 值用引号括起来。

QuoteCharacter

字符串

引号字符。默认值:"

BetweenConfiguration 属性

下表描述了 BETWEEN 处理属性:

属性

Mode

说明

LowBoundKey

QUERY_PARAMS

下限的键模板。支持 {FIELD} 占位符。

HighBoundKey

QUERY_PARAMS

上限的键模板。省略则删除上限。

Template

FILTER_STRING

包含 {FIELD}{LOW}{HIGH} 占位符的模板。

支持的运算符

筛选条件谓词中支持以下运算符:EQUAL_TOGREATER_THANLESS_THANGREATER_THAN_OR_EQUAL_TOLESS_THAN_OR_EQUAL_TONOT_EQUAL_TOCONTAINSBETWEENANDOR

字段级覆盖

您可以使用架构定义中的 FilterOverrides 配置每个字段的筛选行为。以下覆盖属性可用:

  • FieldName:覆盖筛选条件输出中使用的字段名称。

  • OperatorMappings:字段级运算符覆盖。

  • BetweenConfiguration:每个字段的 BETWEEN 覆盖。

  • DateTimeFormat:每个字段的 DateTime 格式覆盖。

示例:QUERY_PARAMS 模式

以下示例显示了 REST API 的 FilterConfiguration,该 API 使用括号符号表示操作符的查询参数和 epoch 时间戳表示日期值:

"FilterConfiguration": { "FilterMode": "QUERY_PARAMS", "OperatorMappings": { "EQUAL_TO": "{FIELD}", "GREATER_THAN_OR_EQUAL_TO": "{FIELD}[gte]", "LESS_THAN_OR_EQUAL_TO": "{FIELD}[lte]" }, "DateTimeFormat": "EPOCH_SECONDS", "BetweenConfiguration": { "LowBoundKey": "{FIELD}[gte]", "HighBoundKey": "{FIELD}[lte]" } }

使用此配置,created >= 2024-01-01 AND created <= 2024-06-01 的输入筛选条件将生成 URL 查询字符串:?created[gte]=1704067200&created[lte]=1717200000

示例:FILTER_STRING 模式

以下示例显示了 REST API 的 FilterConfiguration,该 API 使用带有空格填充运算符的单个筛选字符串参数:

"FilterConfiguration": { "FilterMode": "FILTER_STRING", "OperatorMappings": { "EQUAL_TO": " eq ", "GREATER_THAN": " gt ", "GREATER_THAN_OR_EQUAL_TO": " ge ", "LESS_THAN": " lt ", "AND": " and ", "OR": " or " }, "DateTimeFormat": "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", "BetweenConfiguration": { "Template": "{FIELD} ge {LOW} and {FIELD} le {HIGH}" }, "FilterStringConfiguration": { "FilterStringKey": "search", "QuoteStringValues": true, "QuoteCharacter": "\"" } }

使用此配置,status = "ACTIVE" AND lastUpdated > 2024-01-01T00:00:00.000Z 的输入筛选条件将生成 URL 查询字符串:?search=status eq "ACTIVE" and lastUpdated gt "2024-01-01T00:00:00.000Z"

对查询进行分区

您可以将数据读取拆分为跨 Spark 个工作线程的并行分区,以提高吞吐量。REST API 连接器支持基于字段的分区,它根据指定的字段将数据划分为范围。

Spark 作业参数

以下连接选项控制分区行为:

  • PARTITION_FIELD:要分区的字段。必须在架构中标记为 IsPartitionable: true

  • LOWER_BOUND:分区范围的包含下限。

  • UPPER_BOUND:分区范围的上限。中间分区不包含此值。最后一个分区包含此值。

  • NUM_PARTITIONS:并行分区的数量。

基于字段的分区

基于字段的分区根据指定字段将数据划分为多个范围。以下示例基于 lastUpdated 字段对数据进行分区:

rest_read = glueContext.create_dynamic_frame.from_options( connection_type="rest", connection_options={ "connectionName": "connection-name", "ENTITY_NAME": "entity-name", "CONNECTION_TYPE": "REST-connection-type", "PARTITION_FIELD": "lastUpdated", "LOWER_BOUND": "2024-01-01T00:00:00.000Z", "UPPER_BOUND": "2024-12-31T00:00:00.000Z", "NUM_PARTITIONS": "4" } )

lastUpdated 字段上设置 4 个分区时,工作分布如下:

  • 工作线程 1:lastUpdated >= "2024-01-01" AND lastUpdated < "2024-04-01"

  • 工作线程 2:lastUpdated >= "2024-04-01" AND lastUpdated < "2024-07-01"

  • 工作线程 3:lastUpdated >= "2024-07-01" AND lastUpdated < "2024-10-01"

  • 工作线程 4:lastUpdated >= "2024-10-01" AND lastUpdated <= "2024-12-31"

在 RegisterConnectionType 中配置分区支持

要启用基于字段的分区,请在架构定义中将字段标记为可分区:

"Schema": { "lastUpdated": { "FieldDataType": "TIMESTAMP", "IsPartitionable": true } }

以下 FieldDataType 值支持分区:

  • TIMESTAMP:DateTime 分区。将范围拆分为时间窗口。

  • INTEGER:整数分区。将范围拆分为数字窗口。

注意

连接器使用 AND 将分区筛选条件与用户筛选条件组合在一起。如果您的筛选条件包含 LIMIT 子句,则连接器会跳过分区。如果分区生成失败,连接器会回退到单个分区,作业仍会完成。