Python编程向Azure Databricks表插入数组/Map数据失败求助
向Azure Databricks插入数组/Map类型数据的两类问题及解决办法
一、SQLAlchemy方式的错误与修复
原代码
from sqlalchemy import create_engine, Column, String, text, Integer, JSON, BigInteger, Identity, ARRAY from sqlalchemy.orm import sessionmaker, declarative_base payload = { "schema_name": "catalog_schema", "table_name": "policy11", "column_names": ["CustomerId", "FirstName", "LastName", "Email", "DOB", "Gender", "AnnualIncome", "StreetAddress", "City", "State", "Country", "Zip", "CreatedDate", "UpdatedDate"], "choice": {"AnomalyDetection": "1", "BusinessContext": "1", "DQRules": "1", "Standardization": "0"}, "standardization": [] } Base = declarative_base() class DataQuality(Base): __tablename__ = 'data_quality' id = Column(BigInteger, Identity(start=1, increment=1), primary_key=True) schema_name = Column(String, nullable=False) table_name = Column(String, nullable=False) column_names = Column(ARRAY(String)) choice = Column(JSON()) standardization = Column(ARRAY(String)) engine = create_engine( url = f"databricks://token:{access_token}@{server_hostname}?" + f"http_path={http_path}&catalog={catalog}&schema={schema}" ) Session = sessionmaker(bind=engine) session = Session() def insert_data(session, data): try: result = {} result["schema_name"]=data["schema_name"] result["table_name"]=data["table_name"] result["column_names"]=data["column_names"] result["choice"]=str(data["choice"]) result["standardization"]=data.get("standardization", []) # Add and commit the new record dq_config = DataQuality(**result) session.add(dq_config) session.commit() except Exception as e: session.rollback() print(e) finally: session.close() insert_data(session, payload)
错误信息
(builtins.AttributeError) 'DatabricksDialect' object has no attribute '_json_serializer'
修复方案
该错误源于Databricks的SQLAlchemy方言未默认配置JSON序列化器,同时你将choice转为字符串的处理不符合JSON字段的预期,可通过以下两种方式修复:
- 配置JSON序列化器:创建引擎时指定序列化器,同时直接传入字典而非字符串:
import json # 修改引擎创建逻辑 engine = create_engine( url = f"databricks://token:{access_token}@{server_hostname}?" + f"http_path={http_path}&catalog={catalog}&schema={schema}", json_serializer=json.dumps ) # 修改insert_data中的choice赋值 result["choice"] = data["choice"] # 无需转为字符串
- 改用Text类型存储JSON字符串:若不想配置序列化器,可将
choice字段改为Text类型,存入标准JSON字符串:
from sqlalchemy import Text # 修改DataQuality类的choice字段 class DataQuality(Base): __tablename__ = 'data_quality' # 其他字段保持不变 choice = Column(Text) # 替换原JSON()类型 # insert_data中用json.dumps生成标准JSON字符串 result["choice"] = json.dumps(data["choice"])
二、databricks-sql-connector方式的错误与修复
原代码
connection = sql.connect(server_hostname = server_hostname, http_path = http_path, access_token = access_token) cursor = connection.cursor() cursor.execute("INSERT INTO dq_config_driven_execution.dq_configuration_test.data_quality (schema_name, table_name, column_names, choice) VALUES ('schema_name', 'table_name', ['CustomerId', 'FirstName'], {'AnomalyDetection': '1'})") cursor.close() connection.close()
错误信息
databricks.sql.exc.ServerOperationError: [PARSE_SYNTAX_ERROR] Syntax error at or near '['. SQLSTATE: 42601 (line 1, pos 159)
修复方案
Databricks SQL不支持直接使用Python风格的[]和{}表示数组与Map,需使用Databricks内置构造函数或参数化查询解决:
- 使用Databricks SQL构造函数:
cursor.execute(""" INSERT INTO dq_config_driven_execution.dq_configuration_test.data_quality (schema_name, table_name, column_names, choice) VALUES ('schema_name', 'table_name', array('CustomerId', 'FirstName'), map('AnomalyDetection', '1')) """)
- 参数化查询(推荐):更安全且适配动态数据,避免语法拼接问题:
sql_query = """ INSERT INTO dq_config_driven_execution.dq_configuration_test.data_quality (schema_name, table_name, column_names, choice) VALUES (?, ?, ?, ?) """ params = ('schema_name', 'table_name', ["CustomerId", "FirstName"], {"AnomalyDetection": "1"}) cursor.execute(sql_query, params)
内容的提问来源于stack exchange,提问作者Chinmaya
相关产品推荐
相关产品推荐

