如何在Databricks中压缩JSON Lines文件并上传至Azure ADLS Gen2容器?
解决方案
分为两种实现方案,你可以根据自己的场景选择:
方案1:基于现有流程最小改动
仅需在原有代码基础上增加gzip压缩逻辑,无需调整整体流程:
- 新增依赖导入,将生成的JSON Lines字符串压缩为gzip字节流
import gzip import io # 原有转换逻辑保留 pandas_df = df.select("*").toPandas() json_lines_data = pandas_df.to_json(orient='records', lines=True) # 新增压缩逻辑 compressed_buffer = io.BytesIO() with gzip.GzipFile(fileobj=compressed_buffer, mode='w') as f: f.write(json_lines_data.encode('utf-8')) compressed_data = compressed_buffer.getvalue()
- 修改上传函数,新增压缩相关的响应头配置
def upload_blob(compressed_data, connection_string, container_name, blob_name): blob_service_client = BlobServiceClient.from_connection_string(connection_string) blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name) # 原有删除逻辑可以替换为overwrite参数,更简洁 blob_client.upload_blob( compressed_data, overwrite=True, content_type='application/json', content_encoding='gzip' )
注意:3GB数据全量拉到Driver端转为pandas对象,很容易导致Driver内存不足OOM,仅建议数据量不超过Driver可用内存的场景使用该方案。
方案2:Databricks原生优化方案(更推荐)
不需要将数据拉到Driver端转pandas,利用Spark分布式能力直接生成单文件、自定义文件名的压缩JSON Lines,性能和稳定性更高,完美解决你需要单文件自定义文件名的需求:
- 先将数据合并为单分区,写入压缩JSON Lines到临时路径
# 替换为自己的ADLS路径 temp_output_path = "abfss://<容器名>@<存储账号>.dfs.core.chinacloudapi.cn/temp_json_output/" target_path = "abfss://<容器名>@<存储账号>.dfs.core.chinacloudapi.cn/<自定义文件名>.json.gz" # 分布式写入压缩JSON Lines df.coalesce(1)\ .write\ .mode("overwrite")\ .option("compression", "gzip")\ .option("lineSep", "\n")\ .json(temp_output_path)
- 重命名临时生成的分片文件为目标文件名,清理临时目录
# 查找生成的压缩分片文件 part_file_path = [f.path for f in dbutils.fs.ls(temp_output_path) if f.name.endswith(".json.gz")][0] # 重命名为自定义文件名 dbutils.fs.mv(part_file_path, target_path) # 清理临时目录 dbutils.fs.rm(temp_output_path, recurse=True)
该方案的优势:
- 全流程分布式计算,不会占用Driver内存,无OOM风险
- 压缩和写入效率远高于Driver端单进程处理
- 无需额外依赖Azure存储SDK,直接用Databricks内置工具即可完成
- 生成的压缩文件符合标准gzip规范,客户端访问时可自动识别解压
内容的提问来源于stack exchange,提问作者WIT
相关产品推荐
相关产品推荐

