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

使用Spark DataFrame读取S3文件时Glue Bookmark失效问题求助

解答你的AWS Glue两个核心问题

一、Glue Bookmark无法增量读取.gz JSON文件的问题

你提到用Spark DataFrame读取.gz文件时,Glue Bookmark完全不起作用,每次都全量扫描——这很正常,因为Spark原生的DataFrame读取API(比如spark.read.json)根本不会和Glue Bookmark机制交互。Bookmark是Glue专属的功能,只和Glue的DynamicFrame以及它的数据源绑定。

所以确实需要改用Glue DynamicFrame来读取数据,这样Bookmark才能追踪已经处理过的.gz文件,实现真正的增量读取。而且Glue对.gz压缩的JSON文件支持很好,不需要额外的解压操作,只要配置正确就行。

修改后的核心代码片段

把你原来的srcDf = spark.read.json("s3://source/*/*")替换成Glue的DynamicFrame读取方式,然后再转换成DataFrame处理(如果你习惯用Spark API做转换的话):

# 改用Glue DynamicFrame读取,开启Bookmark支持
srcDyf = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={
        "path": "s3://source/",
        "recurse": True,  # 递归读取子文件夹
        "groupFiles": "inPartition",  # 可选:按分区分组文件,优化性能
        "enableUpdateCatalog": False  # 如果不需要自动更新Catalog可以设为False
    },
    format="json",
    transformation_ctx="srcDyf"  # 关键:这个参数是Bookmark追踪的标识,必须设置
)

# 转换成Spark DataFrame做后续处理(保留你原来的转换逻辑)
srcDf = srcDyf.toDF()

这里的transformation_ctx是核心——Glue通过这个参数来关联Bookmark的状态,只要这个参数不变,每次运行作业时就会自动跳过已经处理过的文件。另外,Glue会自动识别.gz压缩的文件,不需要额外指定压缩格式,它会自动解压读取。

然后你原来的后续转换和写入逻辑可以保留,不过建议写入时也用Glue的DynamicFrame写入API(你已经在用了),这样整个作业的Bookmark状态会更完整。

二、Glue Crawler针对日期子文件夹生成大量表的问题

当你的S3桶有大量按日期命名的子文件夹(比如s3://source/2024/05/20/、s3://source/2024/05/21/)时,Crawler默认可能会把每个子文件夹当成一个独立的表,这是因为它没有识别出这些是分区目录。解决这个问题有几个方案:

方案1:配置Crawler识别分区键

在创建Crawler时,在“配置输出”步骤里,找到“分区和分类”设置:

  • 勾选“启用分区发现”
  • 指定分区键的格式,比如你的日期文件夹是year=2024/month=05/day=20(如果是这种格式的话),或者如果你的文件夹是纯数字(比如2024/05/20),可以设置分区键为year、month、day,并指定对应的路径模式。

如果你的文件夹路径是s3://source/{year}/{month}/{day}/,可以在Crawler的数据源配置里,设置“包含路径”为s3://source/,然后在分区发现里设置分区键的顺序为year、month、day,Crawler就会把这些子文件夹识别为同一个表的分区,而不是生成多个表。

方案2:调整Crawler的扫描策略

  • 设置Crawler的“扫描范围”为“仅扫描新文件夹”,这样它不会重复扫描旧的文件夹,也不会因为旧文件夹的结构重复生成表。
  • 确保你的所有JSON文件的Schema是一致的——如果不同文件夹下的JSON Schema不一样,Crawler会认为是不同的表,所以要保证数据格式统一。

方案3:手动创建表并加载分区

如果Crawler还是不听话,你可以手动在Glue Data Catalog里创建一个表,定义好Schema和分区键(year、month、day),然后用MSCK REPAIR TABLE命令(或者通过Glue的API)来自动加载所有分区,这样就不会生成大量表了。

完整修改后的作业代码

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
from pyspark.sql.types import IntegerType, TimestampType, LongType
from pyspark.sql.functions import col, year, month, dayofmonth, to_date, from_unixtime
from awsglue.dynamicframe import DynamicFrame

## @params: [JOB_NAME]
args = getResolvedOptions(sys.argv, ['JOB_NAME'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# 改用Glue DynamicFrame读取,支持Bookmark增量读取.gz JSON文件
srcDyf = glueContext.create_dynamic_frame.from_options(
    connection_type="s3",
    connection_options={
        "path": "s3://source/",
        "recurse": True
    },
    format="json",
    transformation_ctx="srcDyf"  # 必须设置这个参数让Bookmark追踪状态
)

# 转换为Spark DataFrame进行处理
srcDf = srcDyf.toDF()

# 保留你原来的分区逻辑
partitionDf = srcDf.withColumn("date_col", to_date(col("timestamp"), 'yyyy-MM-dd')) \
                   .withColumn("year", year(col("date_col"))) \
                   .withColumn("month", month(col("date_col"))) \
                   .withColumn("day", dayofmonth(col("date_col"))) \
                   .repartition(1)

# 转换回DynamicFrame进行写入
dynamicdf = DynamicFrame.fromDF(partitionDf, glueContext, "test_nest")

# 写入Parquet到目标S3桶,保留分区
apilogs = glueContext.write_dynamic_frame.from_options(
    frame = dynamicdf,
    connection_type = "s3",
    connection_options = {
        "path": "s3://destination/",
        "partitionKeys": ["year", "month", "day"]
    },
    format = "glueparquet",
    transformation_ctx = "apilogs"
)

job.commit()

注意事项

  • 确保你的Glue作业已经启用了Bookmark(在作业配置里勾选“启用书签”)。
  • 第一次运行作业时会全量读取所有文件,之后每次运行就只会读取新增的.gz文件了。
  • 如果需要重置Bookmark,可以在作业运行时设置参数--job-bookmark-option reset。

内容的提问来源于stack exchange,提问作者Julie C.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:55:30