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

将DataFrame写入Cassandra时遇CassandraCatalogException:缺失主键列

问题描述

在Databricks中使用PySpark尝试将Azure Blob存储中的CSV文件写入Cassandra表:

  • 手动创建或从Blob存储读取的小DataFrame(1-2行)均可正常写入
  • 写入约5万行的大DataFrame时,执行写入命令的单元格报错
写入用PySpark命令
(
 df_csv
  .write.format("org.apache.spark.sql.cassandra")
  .mode('append')
  .options(table=<表名>, keyspace=<键空间>)
  .save()
)
报错信息

com.datastax.spark.connector.datasource.CassandraCatalogException: Attempting to write to C* Table but missing

环境与已做排查
  • 使用的Cassandra驱动:com.datastax.spark:spark-cassandra-connector-assembly_2.12:3.2.0
  • 怀疑报错信息不完整,但无法确认
  • 已核对DataFrame与Cassandra表的列名,二者完全一致
测试代码(小数据写入正常)
from pyspark.sql import Row

department1 = Row(id='1', name='Computer Science')
department2 = Row(id='2', name='Mechanical Engineering')
departments = [department1, department2]

df1 = spark.createDataFrame(departments)

# Spark配置
spark.conf.set("spark.cassandra.connection.host", <主机IP>)
spark.conf.set("spark.cassandra.auth.username", <用户名>)
spark.conf.set("spark.cassandra.auth.password", <密码>)

# 创建Cassandra表
cassandra_session.execute('create table transactions(id int PRIMARY KEY, NAME varchar)')

# 写入小DataFrame
df1.write.format("org.apache.spark.sql.cassandra").mode('append').options(table="transactions", keyspace=<键空间>).save()
解决方案建议
  1. 补全Spark Cassandra连接配置
    大批次写入时需补充连接超时、数据中心指定等参数,避免连接中断:

    spark.conf.set("spark.cassandra.connection.connect_timeout_ms", "10000")
    spark.conf.set("spark.cassandra.connection.read_timeout_ms", "20000")
    spark.conf.set("spark.cassandra.connection.local_dc", "<你的数据中心名称>")
    
  2. 严格核对数据类型匹配
    列名一致但数据类型不匹配会导致批量写入失败,尤其是大批次数据中存在隐式类型转换的情况。执行df_csv.printSchema(),再对比Cassandra表结构:

    DESCRIBE <键空间>.<表名>;
    
  3. 调整批量写入参数
    调整连接器的批量参数,避免超过Cassandra的写入阈值:

    spark.conf.set("spark.cassandra.output.batch.size.rows", "1000")
    spark.conf.set("spark.cassandra.output.concurrent.writes", "10")
    spark.conf.set("spark.cassandra.output.throughput_mb_per_sec", "10")
    
  4. 获取完整报错日志
    在Databricks中点击单元格报错旁的View Full Logs,或查看集群Driver日志,获取完整的异常堆栈,确认具体缺失的内容(通常是主键列、必填字段或权限问题)。

  5. 验证Cassandra账号权限
    确保Spark使用的Cassandra账号拥有目标表的WRITE权限,大批次写入时权限验证逻辑更严格:

    GRANT WRITE ON <键空间>.<表名> TO <用户名>;
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 03:15:49