如何从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接口访问Databrickspyspark:可选,用于复杂数据场景下的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": "数据删除完成"})
注意事项
- 权限验证:确保你的PAT拥有对应Databricks资源的操作权限,可在Databricks工作区的权限设置中调整
- 敏感信息管理:生产环境绝对不要硬编码令牌、URL等敏感信息,用环境变量或Django的
secrets管理 - 性能优化:如果是高频CRUD场景,优先使用Databricks SQL Warehouse而非集群,延迟更低
- Delta特性利用:Delta Lake支持ACID事务,不用担心CRUD操作的原子性问题;还可利用时间旅行特性恢复历史数据
内容的提问来源于stack exchange,提问作者penchalaiah narakatla
相关产品推荐
相关产品推荐

