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

Databricks插入DataFrame至SQL表时重复列异常排查

解决Databricks插入表时的重复列nr_cnae_prin报错

排查与解决步骤

1. 替换SELECT *为显式列名

Spark SQL默认对列名大小写不敏感,若临时视图tmpBcnView中存在大小写/空格差异的同名列(比如nr_cnae_prin和NR_CNAE_PRIN),SELECT *会将其识别为重复列。直接指定列名可避免该问题:

INSERT INTO TABLE db.tb_jul_bcn  
SELECT nr_cpf_cnpj, tp_pess, am_bacen, cd_moda, cd_sub_moda, vl_bacen, clivenc, vl_envio, nm_pess_empr, nr_cnae_prin 
FROM tmpBcnView

2. 检查目标表Schema

确认目标表db.tb_jul_bcn是否存在重复列定义,或列映射冲突:

DESCRIBE TABLE db.tb_jul_bcn;

对比DataFrame Schema,确保列名、数量、类型完全匹配。

3. 排查视图生成过程

若tmpBcnView由JOIN等操作生成,可能在关联时引入了重复列(比如关联的两张表都包含nr_cnae_prin)。重新生成视图时显式指定列并去重:

# 示例:JOIN场景下明确指定列来源
df = df1.join(df2, on="nr_cpf_cnpj", how="inner").select(
    df1.nr_cpf_cnpj,
    df1.tp_pess,
    df1.am_bacen,
    df1.cd_moda,
    df1.cd_sub_moda,
    df1.vl_bacen,
    df1.clivenc,
    df1.vl_envio,
    df1.nm_pess_empr,
    df1.nr_cnae_prin  # 明确选取其中一张表的列
)
df.createOrReplaceTempView("tmpBcnView")

原始问题详情

DataFrame Schema

|-- nr_cpf_cnpj: string (nullable = true)
 |-- tp_pess: string (nullable = true)
 |-- am_bacen: long (nullable = true)
 |-- cd_moda: long (nullable = true)
 |-- cd_sub_moda: long (nullable = true)
 |-- vl_bacen: decimal(29,2) (nullable = true)
 |-- clivenc: string (nullable = true)
 |-- vl_envio: decimal(28,2) (nullable = true)
 |-- nm_pess_empr: string (nullable = true)
 |-- nr_cnae_prin: long (nullable = true)

执行代码

spark.sql("INSERT INTO TABLE db.tb_jul_bcn  SELECT * FROM tmpBcnView")

报错信息

AnalysisException: Found duplicate column(s) in the data to save: nr_cnae_prin
---------------------------------------------------------------------------
AnalysisException                         Traceback (most recent call last)
<command-2987275027841731> in <cell line: 1>()
----> 1 spark.sql("INSERT INTO TABLE db.tb_jul_bcn  SELECT * FROM tmpBcnView")
/databricks/spark/python/pyspark/instrumentation_utils.py in wrapper(*args, **kwargs)
     46             start = time.perf_counter()
     47             try:
---> 48                 res = func(*args, **kwargs)
     49                 logger.log_success(
     50                     module_name, class_name, function_name, time.perf_counter() - start, signature
/databricks/spark/python/pyspark/sql/session.py in sql(self, sqlQuery, **kwargs)
   1117             sqlQuery = formatter.format(sqlQuery, **kwargs)
   1118         try:
---> 1119             return DataFrame(self._jsparkSession.sql(sqlQuery), self)
   1120         finally:
   1121             if len(kwargs) > 0:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 09:12:50