IoT场景下如何用AWS Glue/Athena合并小Parquet文件提升查询性能
可行方案:基于AWS无服务器服务合并小Parquet文件
当然可以实现,以下是几种适配你业务场景的无服务器方案,能按需合并小Parquet文件且保留原桶查询能力:
1. AWS Glue 无服务器ETL作业(推荐)
这是最灵活可控的方案,完全基于无服务器架构,能精准控制合并逻辑。
- 核心步骤:
- 创建Glue无服务器作业,用PySpark编写合并逻辑:
- 读取指定
tablename路径下的所有小Parquet文件(可通过Glue数据目录直接关联表,或直接指定S3路径) - 根据目标文件大小(比如100MB/个)重新分区数据,实现合并
- 将合并后的文件先写入临时路径,再移动到原输出路径(避免覆盖Lambda新生成的文件)
- 可选:过滤并删除旧的小文件(比如只删除24小时前生成的,跳过刚上传的文件)
- 读取指定
- 关键代码示例:
from awsglue.context import GlueContext from pyspark.context import SparkContext import boto3 from datetime import datetime, timedelta sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session s3 = boto3.client('s3') bucket = "outputbucket" target_table_prefix = "*/tablename/" temp_prefix = "temp/tablename/" target_file_size_mb = 100 # 读取目标路径下的所有Parquet数据 df = spark.read.parquet(f"s3://{bucket}/{target_table_prefix}") # 估算分区数(按1GB总数据生成10个100MB文件为例) total_data_mb = 1000 num_partitions = int(total_data_mb / target_file_size_mb) df = df.repartition(num_partitions) # 写入临时路径 df.write.mode("overwrite").parquet(f"s3://{bucket}/{temp_prefix}") # 将临时文件移动到原路径,重命名避免冲突 temp_objects = s3.list_objects_v2(Bucket=bucket, Prefix=temp_prefix) for obj in temp_objects.get('Contents', []): temp_key = obj['Key'] # 生成带merged前缀的新文件名 new_key = f"{target_table_prefix.rstrip('/')}/merged_{obj['Key'].split('/')[-1]}" s3.copy_object(Bucket=bucket, CopySource=f"{bucket}/{temp_key}", Key=new_key) s3.delete_object(Bucket=bucket, Key=temp_key) # 删除24小时前的旧小文件 old_objects = s3.list_objects_v2(Bucket=bucket, Prefix=target_table_prefix) cutoff_time = datetime.now() - timedelta(hours=24) for obj in old_objects.get('Contents', []): if obj['LastModified'].replace(tzinfo=None) < cutoff_time and not obj['Key'].startswith(f"{target_table_prefix.rstrip('/')}/merged_"): s3.delete_object(Bucket=bucket, Key=obj['Key']) - 触发方式:
- 按需手动触发作业
- 通过EventBridge定时触发(比如每日凌晨低峰期执行)
- 基于S3 Inventory + Lambda检测文件数量,当达到阈值(比如某路径下超过1000个小文件)时自动触发Glue作业
- 创建Glue无服务器作业,用PySpark编写合并逻辑:
2. Athena CTAS 语句(无代码方案)
适合简单合并场景,无需编写代码,直接用SQL完成。
- 核心步骤:
- 确保Glue数据目录中已存在对应表的定义
- 执行CTAS语句将数据写入临时路径,再移动到原输出路径:
CREATE TABLE my_table_merged WITH ( format = 'PARQUET', write_compression = 'SNAPPY', external_location = 's3://outputbucket/temp/tablename/' ) AS SELECT * FROM my_table; - 用S3 API将临时路径的合并文件移动到原
*/tablename/路径下 - 可选:删除旧小文件,无需更新Glue表定义(Athena会自动扫描路径下所有文件)
- 注意事项:Athena会自动根据数据量生成合适大小的文件,但无法精确控制文件数量,适合对合并粒度要求不高的场景。
3. Lambda + S3 Batch Operations 方案
适合中小规模文件合并,无需Spark环境,轻量灵活。
- 核心步骤:
- 开启S3 Inventory获取指定路径下的小文件列表
- 编写Lambda函数,用pandas或pyarrow库读取一批小文件并合并成大Parquet文件
- 通过S3 Batch Operations批量触发Lambda处理分组后的文件
- 将合并后的文件写入原路径,删除原小文件
- 注意事项:Lambda有内存和执行时间限制,单批处理的文件数量不宜过多,大规模合并建议优先选择Glue。
关键注意事项
- 避免干扰新数据:合并过程中务必先写入临时路径再移动到原路径,同时通过文件修改时间过滤,不删除最近生成的文件
- 查询兼容性:合并后的Parquet文件与原文件格式一致,Athena可以直接查询混合的新旧文件,无需额外配置
- 成本优化:按需触发(比如文件数达标时)比定时触发更节省成本,Glue无服务器作业按计算时长收费,Athena按扫描数据量收费
内容的提问来源于stack exchange,提问作者mfcss
相关产品推荐
相关产品推荐

