使用Spark DataFrame读取S3文件时Glue Bookmark失效问题求助
一、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.

