如何在AWS Glue作业中使用bucketed_by?求实现数据分区分桶
在AWS Glue作业中实现分桶(bucketed_by)写入
要在Glue作业中直接实现分桶写入,你需要借助Spark原生的DataFrame API(Glue的DynamicFrame Sink目前不直接支持分桶配置)。以下是修改后的完整代码:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) # 读取源CSV数据 S3bucket_node1 = glueContext.create_dynamic_frame.from_options( format_options={"withHeader": True, "separator": ",", "optimizePerformance": False}, connection_type="s3", format="csv", connection_options={"paths": ["s3://source-files-1029/sample_file.csv"]}, transformation_ctx="S3bucket_node1", ) # 字段映射转换 ApplyMapping_node2 = ApplyMapping.apply( frame=S3bucket_node1, mappings=[ ("client_id", "string", "client_id", "string"), ("product_id", "string", "product_id", "string"), ], transformation_ctx="ApplyMapping_node2", ) # 关键步骤1:将Glue DynamicFrame转为Spark DataFrame df = ApplyMapping_node2.toDF() # 关键步骤2:配置分桶+分区并写入数据 # 按product_id分10桶(可根据数据量调整),同时按client_id分区 df.write \ .mode("overwrite") \ .partitionBy("client_id") \ .bucketBy(10, "product_id") \ .format("parquet") \ .option("compression", "snappy") \ .saveAsTable("db1029.tbl-partition-bucket") job.commit()
核心改动说明
- 转换为Spark DataFrame:调用
toDF()方法切换到Spark原生API,才能使用分桶配置能力。 - 分桶配置:通过
bucketBy(桶数, 分桶字段)指定分桶规则,示例中按product_id划分10个桶,你可根据数据规模、查询模式调整参数。 - 分区+分桶共存:保留原有的
partitionBy("client_id")逻辑,分区和分桶可同时生效,优化查询性能。 - 自动同步Glue Catalog:使用
saveAsTable写入时,会自动将数据存入目标表对应的S3路径,同时更新Glue Catalog中的表元数据(包括分桶、分区信息),无需额外手动维护。
注意事项
- 权限配置:确保Glue作业的IAM角色拥有源S3读取、目标S3写入、Glue Catalog表修改的权限。
- 写入模式:示例用
mode("overwrite")覆盖数据,如需追加可改为mode("append"),但分桶表追加需注意数据一致性。 - 桶数选择:建议参考集群核心数倍数或常见查询的过滤维度设置,避免桶数过多/过少影响性能。
内容的提问来源于stack exchange,提问作者Milorad Krstevski
相关产品推荐
相关产品推荐

