You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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中的表元数据(包括分桶、分区信息),无需额外手动维护。

注意事项

  1. 权限配置:确保Glue作业的IAM角色拥有源S3读取、目标S3写入、Glue Catalog表修改的权限。
  2. 写入模式:示例用mode("overwrite")覆盖数据,如需追加可改为mode("append"),但分桶表追加需注意数据一致性。
  3. 桶数选择:建议参考集群核心数倍数或常见查询的过滤维度设置,避免桶数过多/过少影响性能。

内容的提问来源于stack exchange,提问作者Milorad Krstevski

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 09:50:25