从 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]=1717200000FILTER_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 注册中配置筛选行为。下表描述了可用的属性:
属性 |
类型 |
说明 |
|---|---|---|
|
字符串 |
必填项。 |
|
Map |
将逻辑运算符映射到特定于 API 的语法。 |
|
字符串 |
DateTime 格式模式或 |
|
布尔值 |
指定是否从输入值中去除周围的引号。默认值: |
|
对象 |
默认 BETWEEN 处理配置。 |
|
对象 |
特定于 |
FilterStringConfiguration properties
下表描述了 FILTER_STRING 模式的属性:
属性 |
类型 |
说明 |
|---|---|---|
|
字符串 |
必需。查询参数键(例如,“search”或“filter”)。 |
|
布尔值 |
指定是否将 String 和 DateTime 值用引号括起来。 |
|
字符串 |
引号字符。默认值: |
BetweenConfiguration 属性
下表描述了 BETWEEN 处理属性:
属性 |
Mode |
说明 |
|---|---|---|
|
QUERY_PARAMS |
下限的键模板。支持 |
|
QUERY_PARAMS |
上限的键模板。省略则删除上限。 |
|
FILTER_STRING |
包含 |
支持的运算符
筛选条件谓词中支持以下运算符:EQUAL_TO、GREATER_THAN、LESS_THAN、GREATER_THAN_OR_EQUAL_TO、LESS_THAN_OR_EQUAL_TO、NOT_EQUAL_TO、CONTAINS、BETWEEN、AND、OR。
字段级覆盖
您可以使用架构定义中的 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 子句,则连接器会跳过分区。如果分区生成失败,连接器会回退到单个分区,作业仍会完成。