如何将Databricks SQL Connector与Django ORM集成实现CRUD操作?
使用Django ORM连接Databricks的方案
方案一:利用Databricks的PostgreSQL兼容协议(无需自定义适配器)
Databricks SQL支持PostgreSQL兼容的连接方式,直接使用Django官方的PostgreSQL数据库引擎即可,完全无需手写SQL或自定义适配器,能直接复用ORM的所有功能。
配置步骤:
- 安装依赖:
pip install psycopg2-binary # 或psycopg2,根据环境选择
- 修改Django
settings.py中的数据库配置:
import os from urllib.parse import quote_plus DATABASES = { 'default': { 'ENGINE': 'django.db.backends.postgresql', 'NAME': f"{os.getenv('DATABRICKS_CATALOG', 'hive_metastore')}.{os.getenv('DATABRICKS_SCHEMA', 'default')}", 'USER': 'token', 'PASSWORD': os.getenv('DATABRICKS_ACCESS_TOKEN'), 'HOST': os.getenv('DATABRICKS_HOSTNAME').replace('https://', ''), 'PORT': 443, 'OPTIONS': { 'sslmode': 'require', 'connect_timeout': 10, }, } }
关键配置说明:
NAME:必须指定**目录(Catalog)+模式(Schema)**的组合格式,比如hive_metastore.defaultUSER:固定填token,密码字段传入你的Databricks访问令牌HOST:需要去掉开头的https://前缀,例如把https://xxx.cloud.databricks.com改为xxx.cloud.databricks.com
配置完成后,即可正常使用Django ORM的全部功能:模型创建、CRUD操作、数据迁移等,完全不需要编写原生SQL。
方案二:基于databricks-sql-connector自定义Django数据库引擎(进阶)
如果必须使用databricks-sql-connector包实现连接,需要编写自定义数据库引擎适配器,让Django ORM适配Databricks的SQL语法与连接逻辑。
核心步骤:
- 安装依赖:
pip install databricks-sql-connector django
- 项目中创建自定义引擎目录:
在你的Django项目根目录下新建databricks_engine文件夹,结构如下:
your_project/ ├── databricks_engine/ │ ├── __init__.py │ ├── base.py │ ├── client.py │ ├── creation.py │ ├── operations.py │ └── schema.py
- 实现核心连接类(
base.py):
from django.db.backends.base.base import BaseDatabaseWrapper from .client import DatabricksDatabaseClient from .creation import DatabricksDatabaseCreation from .operations import DatabricksDatabaseOperations from .schema import DatabricksDatabaseSchemaEditor import databricks.sql as sql class DatabaseWrapper(BaseDatabaseWrapper): vendor = 'databricks' display_name = 'Databricks SQL' def get_connection_params(self): return { 'server_hostname': self.settings_dict['HOSTNAME'], 'http_path': self.settings_dict['HTTP_PATH'], 'access_token': self.settings_dict['ACCESS_TOKEN'], } def get_new_connection(self, conn_params): return sql.connect(**conn_params) def init_connection_state(self): # 可在此初始化连接状态,比如设置时区 pass def create_cursor(self, name=None): return self.connection.cursor() # 绑定Django数据库引擎的配套组件 client_class = DatabricksDatabaseClient creation_class = DatabricksDatabaseCreation ops_class = DatabricksDatabaseOperations schema_editor_class = DatabricksDatabaseSchemaEditor
- 实现配套辅助类(示例以
client.py为例):
from django.db.backends.base.client import BaseDatabaseClient class DatabricksDatabaseClient(BaseDatabaseClient): def runshell(self): # Databricks SQL不支持交互式shell,直接抛出异常 raise NotImplementedError("Databricks SQL does not support interactive shell.")
其余类(creation.py、operations.py、schema.py)需继承Django对应的基类,并重写Databricks不兼容的方法:
operations.py:重写date_extract_sql、date_trunc_sql等方法,适配Databricks的日期函数语法schema.py:重写create_model、alter_db_table等方法,处理Databricks的目录/模式结构
- 修改Django配置:
DATABASES = { 'default': { 'ENGINE': 'your_project.databricks_engine', 'HOSTNAME': os.getenv('DATABRICKS_HOSTNAME'), 'HTTP_PATH': os.getenv('DATABRICKS_HTTP_PATH'), 'ACCESS_TOKEN': os.getenv('DATABRICKS_ACCESS_TOKEN'), 'CATALOG': os.getenv('DATABRICKS_CATALOG', 'hive_metastore'), 'SCHEMA': os.getenv('DATABRICKS_SCHEMA', 'default'), } }
注意事项:
- 自定义引擎需要适配大量Databricks与标准SQL的语法差异,开发成本较高,稳定性不如方案一
- 模型迁移时需注意Databricks不支持的SQL特性,需在
operations.py中做语法替换
验证连接有效性
配置完成后,运行以下代码验证连接是否正常:
from django.db import connections conn = connections['default'] conn.ensure_connection() cursor = conn.cursor() cursor.execute("SELECT 1") print(cursor.fetchone()) # 输出(1,)表示连接成功
内容的提问来源于stack exchange,提问作者Muhammad Ehtasham
相关产品推荐
相关产品推荐

