AWS Glue PySpark:如何通过查询从DynamoDB加载指定DataFrame数据
解决DynamoDB全量加载过滤慢的问题
当然可以通过在数据源端指定查询条件,只加载符合要求的少量记录,不用全量拉取后再过滤,这能大幅提升处理效率。具体做法如下:
利用DynamoDB原生查询能力在数据源端过滤
- 如果表有主键(Partition Key),优先用
Query操作(只扫描指定分区,效率远高于全表Scan)。在Glue中创建Dynamic Frame时,通过connection_options传入查询条件:import sys from awsglue.context import GlueContext from awsglue.job import Job from pyspark.context import SparkContext sc = SparkContext() glueContext = GlueContext(sc) job = Job(glueContext) dyf = glueContext.create_dynamic_frame.from_options( connection_type="dynamodb", connection_options={ "dynamodb.tableName": "你的表名", "dynamodb.expression": "partitionKey = :pk AND sortKey BETWEEN :sk_start AND :sk_end", "dynamodb.expression.attributes": {":pk": "user_001", ":sk_start": "2024-01-01", ":sk_end": "2024-01-31"} } ) # 后续处理逻辑 job.commit() - 如果必须用全表扫描(比如没有合适的主键条件),可以通过
dynamodb.filter.expression设置过滤规则,让DynamoDB在扫描时直接剔除不符合条件的记录,减少返回的数据量。
- 如果表有主键(Partition Key),优先用
绝对不要先全量加载再过滤
- 别先把800万条数据全部拉到Glue的Dynamic Frame里,再用
filter方法处理。这种方式会先传输所有数据,再在Spark端过滤,耗时会非常长。一定要把过滤逻辑推到DynamoDB端执行。
- 别先把800万条数据全部拉到Glue的Dynamic Frame里,再用
额外优化点
- 给非主键查询字段创建全局二级索引(GSI),这样可以用Query操作针对GSI查询,避免全表Scan。
- 通过
dynamodb.projection.expression指定只需要的字段,比如只返回id、create_time,减少数据传输量。
内容的提问来源于stack exchange,提问作者Tamil selvan
相关产品推荐
相关产品推荐

