如何配置对接两个ADLS Gen1 Parquet数据源的Django后端
Django对接双ADLS存储读写Parquet实现方案
不要把ADLS存储配置到Django的DATABASES节点下——Parquet是文件格式,ADLS是文件存储,不属于Django ORM默认支持的关系型数据源,直接用官方SDK做独立客户端对接逻辑最顺,没有兼容坑。
前置依赖安装
先装必备的第三方包,所有包都经过生产环境验证:
azure-storage-file-datalake:ADLS官方SDK,兼容Gen1/Gen2接口,别用已经停止维护的老版Gen1专属SDKpandas+pyarrow:Parquet序列化/反序列化、内存中数据增改的最优组合,性能远高于纯Python实现的Parquet解析库python-dotenv:加载环境变量里的存储认证信息,避免密钥硬编码到代码里
安装命令直接执行:pip install azure-storage-file-datalake pandas pyarrow python-dotenv
Django配置项
在settings.py里单独加两个ADLS存储的配置块,和默认的数据库配置完全隔离,不管两个存储是同账户不同路径,还是完全独立的ADLS实例,都可以用这套配置:
# settings.py import os from dotenv import load_dotenv load_dotenv() # ADLS Gen1 参考数据源配置(只读权限即可) ADLS_REF_STORAGE = { "account_url": os.getenv("ADLS_REF_ACCOUNT_URL"), "credential": os.getenv("ADLS_REF_SAS_CREDENTIAL"), "file_system_name": os.getenv("ADLS_REF_CONTAINER"), "root_path": "/parquet/reference/form_config" } # ADLS 业务数据存储配置(读写权限) ADLS_BIZ_STORAGE = { "account_url": os.getenv("ADLS_BIZ_ACCOUNT_URL"), "credential": os.getenv("ADLS_BIZ_SAS_CREDENTIAL"), "file_system_name": os.getenv("ADLS_BIZ_CONTAINER"), "root_path": "/parquet/business/form_submit" }
封装ADLS Parquet操作客户端
单独写工具类封装重复的认证、读写逻辑,避免每个视图都重复写流处理代码:
# utils/adls_parquet.py from io import BytesIO import pandas as pd from azure.storage.filedatalake import DataLakeServiceClient from django.conf import settings class ADLSParquetClient: def __init__(self, storage_config): self.service_client = DataLakeServiceClient( account_url=storage_config["account_url"], credential=storage_config["credential"] ) self.fs_client = self.service_client.get_file_system_client(storage_config["file_system_name"]) self.root_path = storage_config["root_path"].strip("/") def _get_full_path(self, relative_path): return f"{self.root_path}/{relative_path.lstrip('/')}" def read_parquet(self, relative_path: str) -> pd.DataFrame: """读取指定相对路径的Parquet文件,返回pandas DataFrame""" file_client = self.fs_client.get_file_client(self._get_full_path(relative_path)) stream = file_client.download_file().readall() return pd.read_parquet(BytesIO(stream), engine="pyarrow") def write_parquet(self, df: pd.DataFrame, relative_path: str, write_mode: str = "overwrite"): """ 写入Parquet文件 *注意:Parquet是列式存储,不支持原地追加行,append模式本质是读全量旧数据合并后重写 """ file_client = self.fs_client.get_file_client(self._get_full_path(relative_path)) if write_mode == "append" and file_client.exists(): old_df = self.read_parquet(relative_path) df = pd.concat([old_df, df], ignore_index=True) buffer = BytesIO() df.to_parquet(buffer, engine="pyarrow", index=False) buffer.seek(0) file_client.upload_data(buffer, overwrite=True) def update_rows(self, relative_path: str, match_filter: dict, update_data: dict) -> int: """按字段匹配条件更新行,返回实际更新的行数""" df = self.read_parquet(relative_path) # 构造行匹配掩码 match_mask = pd.Series(True, index=df.index) for col, val in match_filter.items(): match_mask &= (df[col] == val) affected_rows = match_mask.sum() # 执行更新 for col, val in update_data.items(): df.loc[match_mask, col] = val # 重写文件 self.write_parquet(df, relative_path) return affected_rows
关键提醒:不要花时间找Parquet原地插入、修改单行的方案,Parquet的存储设计从底层就不支持这种操作,所有增改都是内存修改后全量重写。如果单文件数据量超过500M,建议按日期、业务线拆分Parquet文件,避免每次修改都加载过大的数据集拖慢响应。
接口层对接React前端
初始化两个独立的客户端,分别对接参考数据和业务数据存储,写对应接口给前端调用即可:
# views.py from datetime import datetime import pandas as pd from rest_framework.decorators import api_view from rest_framework.response import Response from django.conf import settings from .utils.adls_parquet import ADLSParquetClient # 全局初始化客户端,避免每次请求重复建连 ref_data_client = ADLSParquetClient(settings.ADLS_REF_STORAGE) biz_data_client = ADLSParquetClient(settings.ADLS_BIZ_STORAGE) @api_view(["GET"]) def get_form_options(request): """返回React表单需要的所有参考选项数据""" region_data = ref_data_client.read_parquet("region_list.parquet") category_data = ref_data_client.read_parquet("category_mapping.parquet") return Response({ "region_options": region_data.to_dict("records"), "category_options": category_data.to_dict("records") }) @api_view(["POST"]) def add_submit_row(request): """新增业务数据行""" new_row_df = pd.DataFrame([request.data]) # 按天拆分存储文件,控制单文件大小 save_path = f"submit_{datetime.now().strftime('%Y%m%d')}.parquet" biz_data_client.write_parquet(new_row_df, save_path, write_mode="append") return Response({"code": 0, "msg": "新增成功"}) @api_view(["POST"]) def edit_submit_row(request): """编辑历史业务数据行""" row_id = request.data.pop("row_id") save_path = request.data.pop("storage_path") # 前端列表项携带所属文件路径,减少全目录扫描开销 updated_count = biz_data_client.update_rows( relative_path=save_path, match_filter={"id": row_id}, update_data=request.data ) return Response({"code": 0, "msg": f"成功更新{updated_count}条数据"})
生产环境踩坑提醒
- 权限严格隔离:给参考数据存储的凭证只开只读权限,业务数据存储开读写权限,避免代码bug误改参考数据
- 并发写控制:如果存在多人同时编辑同一份Parquet文件的场景,用Redis加路径级别的分布式锁,避免并发覆盖导致数据丢失
- 备份机制:每次重写Parquet文件前,把旧版本按时间戳备份到单独的归档路径,误操作时可以快速回滚
- 缓存优化:参考数据更新频率极低,直接给
get_form_options接口加10分钟级别的缓存,不用每次请求都读ADLS,响应速度可以从几百毫秒降到几毫秒
内容的提问来源于stack exchange,提问作者Suel Ahmed
相关产品推荐
相关产品推荐

