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

如何从Django应用连接Databricks Delta Tables并执行CRUD操作

Django连接Databricks Delta Tables并实现CRUD操作步骤

前提准备

  • 获取Databricks核心信息:工作区URL、SQL Warehouse/集群的HTTP路径、个人访问令牌(PAT)、目标Delta表所在的数据库名和表名
  • 确保你的PAT拥有目标表的SELECT/INSERT/UPDATE/DELETE权限

步骤1:安装依赖包

执行以下命令安装连接所需的Python库:

pip install databricks-sql-connector pyspark
  • databricks-sql-connector:用于通过标准SQL接口访问Databricks
  • pyspark:可选,用于复杂数据场景下的DataFrame操作

步骤2:配置Django连接参数

在项目的settings.py中添加Databricks配置(建议用环境变量存储敏感信息):

# settings.py
import os

DATABRICKS_CONFIG = {
    "server_hostname": os.getenv("DATABRICKS_HOST"),  # 格式:xxx.cloud.databricks.com
    "http_path": os.getenv("DATABRICKS_HTTP_PATH"),   # SQL Warehouse或集群的HTTP路径
    "access_token": os.getenv("DATABRICKS_TOKEN"),    # 你的个人访问令牌
    "database": os.getenv("DATABRICKS_DB_NAME")       # 目标Delta表所在的数据库
}

步骤3:封装Databricks操作工具类

在你的Django应用下新建databricks_utils.py,封装CRUD操作,避免重复代码:

# databricks_utils.py
from databricks import sql
from django.conf import settings
from pyspark.sql import SparkSession

# -------------------------- 基于SQL的CRUD实现 --------------------------
def get_db_connection():
    """获取Databricks SQL连接"""
    return sql.connect(
        server_hostname=settings.DATABRICKS_CONFIG["server_hostname"],
        http_path=settings.DATABRICKS_CONFIG["http_path"],
        access_token=settings.DATABRICKS_CONFIG["access_token"]
    )

def read_delta_table(table_name):
    """读取Delta表数据,返回字典列表"""
    conn = get_db_connection()
    cursor = conn.cursor()
    cursor.execute(f"SELECT * FROM {settings.DATABRICKS_CONFIG['database']}.{table_name}")
    columns = [col[0] for col in cursor.description]
    data = [dict(zip(columns, row)) for row in cursor.fetchall()]
    cursor.close()
    conn.close()
    return data

def insert_data(table_name, data_dict):
    """插入单条数据,data_dict为键值对(列名: 值)"""
    conn = get_db_connection()
    cursor = conn.cursor()
    cols = ", ".join(data_dict.keys())
    placeholders = ", ".join([f":{k}" for k in data_dict.keys()])
    query = f"INSERT INTO {settings.DATABRICKS_CONFIG['database']}.{table_name} ({cols}) VALUES ({placeholders})"
    cursor.execute(query, data_dict)
    conn.commit()
    cursor.close()
    conn.close()

def update_data(table_name, update_dict, condition):
    """更新数据,condition为SQL WHERE条件语句"""
    conn = get_db_connection()
    cursor = conn.cursor()
    set_clause = ", ".join([f"{k} = :{k}" for k in update_dict.keys()])
    query = f"UPDATE {settings.DATABRICKS_CONFIG['database']}.{table_name} SET {set_clause} WHERE {condition}"
    cursor.execute(query, update_dict)
    conn.commit()
    cursor.close()
    conn.close()

def delete_data(table_name, condition):
    """删除数据,condition为SQL WHERE条件语句"""
    conn = get_db_connection()
    cursor = conn.cursor()
    query = f"DELETE FROM {settings.DATABRICKS_CONFIG['database']}.{table_name} WHERE {condition}"
    cursor.execute(query)
    conn.commit()
    cursor.close()
    conn.close()

# -------------------------- 基于PySpark的操作(可选) --------------------------
def get_spark_session():
    """获取Spark会话,用于复杂数据处理"""
    return SparkSession.builder \
        .appName("Django-Databricks-Integration") \
        .config("spark.databricks.service.server.enabled", "true") \
        .config("spark.databricks.service.host", settings.DATABRICKS_CONFIG["server_hostname"]) \
        .config("spark.databricks.service.token", settings.DATABRICKS_CONFIG["access_token"]) \
        .getOrCreate()

def read_delta_with_spark(table_name):
    """用PySpark读取Delta表,返回Pandas格式字典列表"""
    spark = get_spark_session()
    df = spark.read.format("delta").table(f"{settings.DATABRICKS_CONFIG['database']}.{table_name}")
    return df.toPandas().to_dict("records")

步骤4:在视图中调用工具类实现业务逻辑

以views.py为例,编写接口调用CRUD方法:

# views.py
from django.http import JsonResponse
from .databricks_utils import read_delta_table, insert_data, update_data, delete_data

def fetch_table_data(request):
    """读取Delta表数据接口"""
    table_data = read_delta_table("your_delta_table_name")
    return JsonResponse(table_data, safe=False)

def add_new_record(request):
    """插入新数据接口"""
    sample_data = {"user_id": 101, "user_name": "John Doe", "email": "john@example.com"}
    insert_data("your_delta_table_name", sample_data)
    return JsonResponse({"status": "success", "msg": "数据插入完成"})

def update_existing_record(request):
    """更新数据接口"""
    update_info = {"email": "john.doe@example.com"}
    update_data("your_delta_table_name", update_info, "user_id = 101")
    return JsonResponse({"status": "success", "msg": "数据更新完成"})

def delete_record(request):
    """删除数据接口"""
    delete_data("your_delta_table_name", "user_id = 101")
    return JsonResponse({"status": "success", "msg": "数据删除完成"})

注意事项

  1. 权限验证:确保你的PAT拥有对应Databricks资源的操作权限,可在Databricks工作区的权限设置中调整
  2. 敏感信息管理:生产环境绝对不要硬编码令牌、URL等敏感信息,用环境变量或Django的secrets管理
  3. 性能优化:如果是高频CRUD场景,优先使用Databricks SQL Warehouse而非集群,延迟更低
  4. Delta特性利用:Delta Lake支持ACID事务,不用担心CRUD操作的原子性问题;还可利用时间旅行特性恢复历史数据

内容的提问来源于stack exchange,提问作者penchalaiah narakatla

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:15:02