AWS Glue如何从MySQL数据源或Dynamic DataFrame筛选记录子集
问题原因
你使用的push_down_predicate参数仅适用于Glue Data Catalog中登记的分区表(通常是S3等对象存储上的分区数据集),对JDBC类型的MySQL数据源不生效,所以会返回全量数据。
方法1:数据源端直接限制拉取(推荐,性能最优)
直接在调用from_catalog时通过connection_options传入自定义查询,让过滤逻辑在MySQL端执行,只返回你需要的子集,代码示例:
from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext.getOrCreate() glueContext = GlueContext(sc) dy_f = glueContext.create_dynamic_frame.from_catalog( database="db", table_name="table with 1000 records", connection_options={ # 这里写你需要的过滤查询,也可以加LIMIT限制返回条数 "query": "SELECT * FROM `table with 1000 records` WHERE id IN (1,2,3)" } ) dy_f.count() # 此时会返回3
如果你用的是Glue 4.0及以上版本,也可以开启JDBC谓词下推功能,后续对DynamicFrame/Spark DataFrame的过滤条件会自动下推到MySQL端执行:
dy_f = glueContext.create_dynamic_frame.from_catalog( database="db", table_name="table with 1000 records", additional_options={ "enableQueryPushdown": "true" } ) # 后续的过滤会自动下推到MySQL执行,不需要全量拉取 filtered_dyf = dy_f.filter(f=lambda x: x["id"] in [1,2,3])
方法2:拉取全量后本地过滤(仅适合小表使用)
如果不方便在数据源端做过滤,也可以拉取全量数据后在Glue侧过滤DynamicFrame,有两种常用写法:
- 直接使用DynamicFrame的filter方法
# 先拉全量 dy_f = glueContext.create_dynamic_frame.from_catalog( database="db", table_name="table with 1000 records" ) # 过滤id符合要求的记录 filtered_dyf = dy_f.filter(f=lambda record: record["id"] in [1,2,3]) # 也可以用limit方法直接取前N条 top_10_dyf = dy_f.limit(10)
- 转成Spark DataFrame过滤后转回DynamicFrame(适合复杂过滤逻辑)
from awsglue.dynamicframe import DynamicFrame dy_f = glueContext.create_dynamic_frame.from_catalog( database="db", table_name="table with 1000 records" ) spark_df = dy_f.toDF() filtered_spark_df = spark_df.filter(spark_df.id.isin([1,2,3])) # 转回DynamicFrame即可继续使用Glue的相关能力 filtered_dyf = DynamicFrame.fromDF(filtered_spark_df, glueContext, "filtered_dyf")
内容的提问来源于stack exchange,提问作者Elijah Bernstein-Cooper
相关产品推荐
相关产品推荐

