使用Databricks AutoLoader能否不合并Parquet,实现S3到Azure Blob的1:1复制?
解决S3到Azure Blob的Parquet文件1:1复制问题
你的问题出在Spark Structured Streaming的默认行为上——它会将多个源文件合并到一个微批处理,并且输出时会自动合并小文件以优化存储效率,这和你想要的1:1文件复制逻辑冲突。mergeSchema参数仅用于处理Schema演化,和文件合并完全无关,所以设置它无法解决问题。
方案一:使用批处理实现严格1:1复制(推荐静态数据场景)
如果是处理已存在的静态数据,直接遍历每个源Parquet文件,逐个读取并写入,确保每个源文件对应一个输出文件:
cda_path = f'{planet}/{center}/{deployment}' for table_name in tables: source_dir = f'dbfs:/mnt/gwcp/{cda_path}/{table_name}' # 获取目录下所有Parquet文件 files = dbutils.fs.ls(source_dir) parquet_files = [f.path for f in files if f.path.endswith('.parquet')] blob_output_base = f"/mnt/test_mount_databricks/{planet}/{center}/{deployment}/{table_name}_test" for file_path in parquet_files: # 读取单个Parquet文件 df = spark.read.parquet(file_path) # 获取源文件名,保持输出文件名一致 file_name = file_path.split('/')[-1] temp_output_dir = f"{blob_output_base}/temp_{file_name}" # 用coalesce(1)强制生成单个文件 df.coalesce(1).write.mode("append").parquet(temp_output_dir) # 将生成的part文件重命名为源文件名 generated_files = dbutils.fs.ls(temp_output_dir) part_file = next(f.path for f in generated_files if f.path.startswith(f"{temp_output_dir}/part-")) dbutils.fs.mv(part_file, f"{blob_output_base}/{file_name}") # 删除临时目录 dbutils.fs.rm(temp_output_dir, recurse=True)
方案二:调整流处理参数适配1:1复制(实时新增文件场景)
如果需要监听S3目录的实时新增文件并复制,需要调整流处理的参数,让每个源文件单独触发一个微批,并且禁止输出文件合并:
cda_path = f'{planet}/{center}/{deployment}' streams = list() for table_name in tables: streams.append( (table_name, spark.readStream.format("cloudFiles") .option("cloudFiles.format", "parquet") .option('cloudFiles.schemaLocation', f'dbfs:/FileStore/shared_uploads/checkpints/stream_{table_name}') .option('cloudFiles.schemaEvolutionMode', 'rescue') # 每次触发仅处理1个源文件 .option('cloudFiles.maxFilesPerTrigger', 1) .load(f'dbfs:/mnt/gwcp/{cda_path}/{table_name}/*')) ) for table_name, stream in streams: blob_output_path = f"/mnt/test_mount_databricks/{planet}/{center}/{deployment}/{table_name}_test" stream.writeStream .format("parquet") .outputMode("append") .option("checkpointLocation", f"dbfs:/FileStore/shared_uploads/checkpoints/blob_{table_name}_test") .option("mergeSchema", "false") # 禁用文件合并相关优化 .option("spark.sql.files.maxRecordsPerFile", 1000000) # 设为远大于单个源文件的行数 .option("spark.sql.streaming.fileSink.log.cleanupDelay", "0") .option("spark.sql.streaming.fileSink.log.compactInterval", "1") .start(blob_output_path)
参数说明:
cloudFiles.maxFilesPerTrigger: 限制每次微批处理的文件数量为1,确保每个源文件单独处理spark.sql.files.maxRecordsPerFile: 设置单个输出文件的最大记录数,值要大于单个源文件的行数,保证整个微批的数据写入一个文件spark.sql.streaming.fileSink.log.cleanupDelay/compactInterval: 加快流处理元数据的清理和压缩,避免后续批次合并已生成的文件
注意:流处理方式无法完全保证绝对的1:1(比如极端情况下的批次触发延迟),如果对文件对应关系有严格要求,优先选择批处理方案。
内容的提问来源于stack exchange,提问作者TechLuk
相关产品推荐
相关产品推荐

