如何让AWS Glue Crawler为单个文件中的多Schema JSON数据生成多张表?
如何让AWS Glue Crawler为单文件内的多Schema JSON生成独立表
默认情况下,AWS Glue Crawler会把单个文件里的所有JSON对象当作同一数据集处理,合并所有字段生成一张表——这确实不是你想要的结果。要让Crawler为每种Schema生成独立表,最可靠的方式是先把不同Schema的对象拆分到单独的S3路径下,再让Crawler爬取这些路径。下面是具体步骤:
步骤1:用Glue ETL拆分多Schema文件
首先创建一个Glue ETL Job,用PySpark脚本把原始文件中的不同JSON对象按Schema拆分到不同的S3位置。这里利用每个Schema的特征字段来过滤:
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 # 解析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) # 读取原始JSON Lines文件(你的文件是每行一个JSON对象,符合JSON Lines格式) raw_df = spark.read.json("s3://your-bucket/path/your-multi-schema-file.json") # 根据特征字段拆分不同Schema的数据集 # 匹配包含"0"字段的Schema schema1_df = raw_df.filter(raw_df["0"].isNotNull()) # 匹配包含event和userId字段的Schema schema2_df = raw_df.filter(raw_df["event"].isNotNull() & raw_df["userId"].isNotNull()) # 匹配包含browser.name字段的Schema(注意用反引号转义带点的字段名) schema3_df = raw_df.filter(raw_df["`browser.name`"].isNotNull()) # 将拆分后的数据集写入不同S3路径 schema1_df.write.mode("overwrite").json("s3://your-bucket/processed-data/schema-type-1/") schema2_df.write.mode("overwrite").json("s3://your-bucket/processed-data/schema-type-2/") schema3_df.write.mode("overwrite").json("s3://your-bucket/processed-data/schema-type-3/") job.commit()
你可以根据实际的Schema特征调整过滤条件,比如新增更多分支处理其他Schema类型。如果原始文件是持续生成的,还可以设置Job定时运行,或者用S3事件触发Job自动处理新文件。
步骤2:配置Glue Crawler生成独立表
拆分完成后,创建或修改Glue Crawler:
- 设置Crawler的数据源为拆分后的根路径
s3://your-bucket/processed-data/ - 在Crawler的配置中,确保表生成策略设为
Create tables automatically(默认就是这个) - 可以设置
Table prefix(比如processed_),让生成的表名更清晰(比如processed_schema_type_1)
运行Crawler后,它会识别每个子路径下的同一Schema数据,自动为每个子路径生成独立的表——就像处理多个独立文件一样。
可选优化:按时间分区
如果你的JSON对象都带有time字段,可以在写入拆分数据时按时间分区,这样能大幅提高后续查询的效率:
# 示例:将time字段转换为日期格式,按日期分区 from pyspark.sql.functions import from_unixtime, date_format schema1_df = schema1_df.withColumn("date", date_format(from_unixtime(schema1_df["time"]/1000), "yyyy-MM-dd")) schema1_df.write.mode("overwrite").partitionBy("date").json("s3://your-bucket/processed-data/schema-type-1/")
这样Crawler还会自动识别分区列,进一步优化表的结构。
内容的提问来源于stack exchange,提问作者Arman Grigoryan
相关产品推荐
相关产品推荐

