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

基于FastAPI向Databricks Catalog写入数据的API端点问题

解决Databricks REST API创建表并插入数据的问题

你用的/api/2.0/sql/createTable不是Databricks官方支持的REST API端点,这是请求失败的核心原因。正确的做法是使用SQL语句执行端点来完成表创建和数据插入操作,具体步骤如下:

1. 替换为正确的API端点

Databricks用来执行SQL语句的REST API端点是:

{databricks_host}/api/2.0/sql/statements

2. 调整代码逻辑(构造SQL语句)

该端点通过执行SQL操作数据,你需要将DataFrame的数据转换为可执行的SQL语句,比如先创建目标表(如果不存在),再插入数据。以下是修改后的完整代码示例:

import requests
import pandas as pd

def send_to_dtb_catalog(self, df: pd.DataFrame, table_name: str):
    # 定义目标表的全限定名
    full_table_name = f"my_database.my_schema.{table_name}"
    
    # 生成创建表的SQL(自动推断DataFrame字段类型,可按需调整)
    create_table_cols = []
    dt_map = {
        'object': 'STRING',
        'int64': 'BIGINT',
        'float64': 'DOUBLE',
        'datetime64[ns]': 'TIMESTAMP'
    }
    for col, dtype in df.dtypes.items():
        sql_type = dt_map.get(str(dtype), 'STRING')
        create_table_cols.append(f"`{col}` {sql_type}")
    
    create_table_sql = f"""
    CREATE TABLE IF NOT EXISTS {full_table_name} (
        {', '.join(create_table_cols)}
    )
    """
    
    # 生成插入数据的SQL(处理数据转义,避免SQL注入)
    def escape_value(val):
        if pd.isna(val):
            return 'NULL'
        elif isinstance(val, str):
            return f"'{val.replace('\\', '\\\\').replace("'", "\\'")}'"
        elif isinstance(val, pd.Timestamp):
            return f"'{val.strftime('%Y-%m-%d %H:%M:%S')}'"
        else:
            return str(val)
    
    rows = []
    for _, row in df.iterrows():
        row_vals = [escape_value(val) for val in row.values]
        rows.append(f"({', '.join(row_vals)})")
    
    insert_sql = f"""
    INSERT INTO {full_table_name}
    VALUES {', '.join(rows)}
    """
    
    # 配置请求头
    headers = {
        "Authorization": f"Bearer {self.databricks_token}",
        "Content-Type": "application/json"
    }
    
    # 先执行创建表请求
    create_response = requests.post(
        url=f"{self.databricks_host}/api/2.0/sql/statements",
        headers=headers,
        json={"statement": create_table_sql}
    )
    create_response.raise_for_status()
    
    # 再执行插入数据请求
    insert_response = requests.post(
        url=f"{self.databricks_host}/api/2.0/sql/statements",
        headers=headers,
        json={"statement": insert_sql}
    )
    insert_response.raise_for_status()
    
    return insert_response.json()

关键注意事项

  • 权限验证:确保你的Databricks Token拥有目标Catalog/Schema的CREATE TABLE和INSERT权限
  • 数据类型扩展:示例中的类型映射是基础版,可根据实际需求(比如Decimal、Array类型)扩展dt_map字典
  • 大数据量优化:如果DataFrame行数过多,建议拆分批量插入,或使用COPY INTO语句从云存储导入数据(更适合大规模数据场景)
  • SQL注入防护:示例中的转义处理是基础版,生产环境建议使用Databricks SQL的参数绑定功能进一步优化

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 05:25:29