You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何配置对接两个ADLS Gen1 Parquet数据源的Django后端

Django对接双ADLS存储读写Parquet实现方案

不要把ADLS存储配置到Django的DATABASES节点下——Parquet是文件格式,ADLS是文件存储,不属于Django ORM默认支持的关系型数据源,直接用官方SDK做独立客户端对接逻辑最顺,没有兼容坑。

前置依赖安装

先装必备的第三方包,所有包都经过生产环境验证:

  • azure-storage-file-datalake:ADLS官方SDK,兼容Gen1/Gen2接口,别用已经停止维护的老版Gen1专属SDK
  • pandas + 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.29 21:30:42