基于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
相关产品推荐
相关产品推荐

