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

如何在非Databricks集群的独立节点上,通过Databricks Python连接器将Pandas DataFrame或CSV数据导入Databricks Delta表?

如何通过Databricks Python连接器将Pandas DataFrame/CSV导入Delta表(无需Spark)

好消息!你完全可以通过Databricks Python连接器实现这个需求,而且不需要依赖Spark。下面我给你两种实用的方案,分别适配不同的场景:

方案1:用executemany批量插入Pandas DataFrame数据

如果你的DataFrame数据量不算特别大(比如百万级以内),这种方法简单直接,不需要额外配置云存储权限:

  1. 先确认目标Delta表已存在(如果没有,用连接器执行CREATE TABLE语句创建)
  2. 将DataFrame转换成元组列表,适配连接器的批量插入格式
  3. 构造带占位符的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表:

  1. 先确保Databricks SQL仓库有权限访问ADLS Gen2(可以通过SAS Token、Azure AD Passthrough或者服务主体配置)
  2. 构造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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 18:17:42