如何在非Databricks集群的独立节点上,通过Databricks Python连接器将Pandas DataFrame或CSV数据导入Databricks Delta表?
如何通过Databricks Python连接器将Pandas DataFrame/CSV导入Delta表(无需Spark)
好消息!你完全可以通过Databricks Python连接器实现这个需求,而且不需要依赖Spark。下面我给你两种实用的方案,分别适配不同的场景:
方案1:用executemany批量插入Pandas DataFrame数据
如果你的DataFrame数据量不算特别大(比如百万级以内),这种方法简单直接,不需要额外配置云存储权限:
- 先确认目标Delta表已存在(如果没有,用连接器执行
CREATE TABLE语句创建) - 将DataFrame转换成元组列表,适配连接器的批量插入格式
- 构造带占位符的
INSERT语句,调用executemany完成批量写入
示例代码:
from databricks import sql import pandas as pd # 假设你已经把ADLS上的CSV读取成了DataFrame df = pd.read_csv("your_local_or_adls_mounted_path") # 建立Databricks连接 conn = sql.connect( server_hostname=self.server_name, http_path=self.http_path, access_token=self.access_token ) try: with conn.cursor() as cursor: # 1. 创建目标表(根据你的实际Schema调整字段和类型) create_table_qry = """ CREATE TABLE IF NOT EXISTS your_target_delta_table ( col1 STRING, col2 INT, col3 DOUBLE ) USING DELTA """ cursor.execute(create_table_qry) # 2. 构造INSERT语句,用%s作为通用占位符 insert_qry = """ INSERT INTO your_target_delta_table (col1, col2, col3) VALUES (%s, %s, %s) """ # 3. 把DataFrame转成元组列表 data_tuples = [tuple(row) for row in df.itertuples(index=False, name=None)] # 4. 批量插入数据 cursor.executemany(insert_qry, data_tuples) # 显式提交事务(更稳妥) conn.commit() print("数据插入成功!") except Exception as e: print(f"出错了:{str(e)}") conn.rollback() finally: conn.close()
注意事项:
- 占位符统一用
%s,不管字段类型是什么,连接器会自动处理类型转换 - 确保DataFrame的列顺序和
INSERT语句中的列顺序完全对应 - 如果数据量超过千万级,这种方法效率会下降,建议用方案2
方案2:用COPY INTO直接从ADLS Gen2导入CSV(高效首选)
既然你的CSV已经存在于ADLS Gen2,用COPY INTO是最高效的方式——不需要把数据拉到本地节点,直接让Databricks从云存储读取并写入Delta表:
- 先确保Databricks SQL仓库有权限访问ADLS Gen2(可以通过SAS Token、Azure AD Passthrough或者服务主体配置)
- 构造
COPY INTO语句,指定ADLS路径、文件格式和认证方式,执行即可
示例代码:
from databricks import sql conn = sql.connect( server_hostname=self.server_name, http_path=self.http_path, access_token=self.access_token ) try: with conn.cursor() as cursor: # 1. 创建目标Delta表(如果不存在) create_table_qry = """ CREATE TABLE IF NOT EXISTS your_target_delta_table ( col1 STRING, col2 INT, col3 DOUBLE ) USING DELTA """ cursor.execute(create_table_qry) # 2. 构造COPY INTO语句,这里用SAS Token认证(替换成你的实际信息) copy_into_qry = """ COPY INTO your_target_delta_table FROM 'abfss://your-container@your-storage-account.dfs.core.windows.net/path/to/your/file.csv' WITH ( FORMAT = 'CSV', HEADER = 'true', # 如果CSV文件有表头就设为true CREDENTIAL = 'your-full-sas-token-here' ) """ # 执行导入 cursor.execute(copy_into_qry) conn.commit() print("CSV数据导入成功!") except Exception as e: print(f"出错了:{str(e)}") conn.rollback() finally: conn.close()
认证方式说明:
- SAS Token:直接把完整的SAS字符串填入
CREDENTIAL参数(用单引号包裹) - Azure AD Passthrough:不需要指定
CREDENTIAL,只要Databricks的服务主体有ADLS的读取权限 - 服务主体:可以在Databricks中创建Secret,然后通过
SECRET(scoped_name, secret_key)引用认证信息
常见问题排查
- 权限错误:检查ADLS的访问控制列表(ACL),确保Databricks SQL仓库有权限读取目标文件
- 类型不匹配:确保目标Delta表的Schema和CSV/DataFrame的字段类型完全一致
- 批量插入超时:可以把DataFrame拆分成多个小批次分多次插入,或者在
connect方法中指定timeout参数调整超时时间
内容的提问来源于stack exchange,提问作者shankar
相关产品推荐
相关产品推荐

