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

Spark+Kafka流写入Cassandra报错:缺失主键列[col1,col2,col3]

Spark Streaming写入Cassandra报错:缺失主键列

运行环境

kafka ----ReadStream----> local ----WriteStream----> cassandra

源代码部署在本地,kafka、本地节点、WriteStream目标节点IP各不相同。

Cassandra表结构

col1 | col2 | col3 | col4 | col5 | col6 | col7

注:根据报错信息,col1、col2、col3为该表的主键列

Spark读取Kafka流后的DataFrame结构

root
|-- key: binary (nullable = true)
|-- value: binary (nullable = true)
|-- topic: string (nullable = true)
|-- partition: integer (nullable = true)
|-- offset: long (nullable = true)
|-- timestamp: timestamp (nullable = true)
|-- timestampType: integer (nullable = true)

运行命令

spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2,com.datastax.spark:spark-cassandra-connector_2.12:3.2.0,com.github.jnr:jnr-posix:3.1.15 \
--conf spark.cassandra.connection.host={cassandraIP} \
--conf spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions \
test.py

注:原命令中com.datastax.spark:spark.cassandra.connectiohost为拼写错误,已修正为正确配置项spark.cassandra.connection.host

源代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col
from pyspark.sql.types import StructType, StructField, StringType, IntegerType  # 根据实际数据类型调整

# 连接Spark集群
spark = SparkSession.builder \
    .master("spark://{SparkMasterIP}:7077") \
    .appName("Spark_Streaming+kafka+cassandra") \
    .config('spark.cassandra.connection.host', '{cassandraIP}') \
    .config('spark.cassandra.connection.port', '9042') \
    .getOrCreate()

# 读取Kafka流
df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "{KafkaIP}:9092") \
    .option('startingOffsets','earliest') \
    .option('failOnDataLoss','False') \
    .option("subscribe", "{Topic}") \
    .load()

df.printSchema()

# 核心步骤:解析Kafka的value字段为Cassandra表对应的结构化数据
# 请根据Kafka消息的实际JSON结构调整Schema字段和类型
kafka_value_schema = StructType([
    StructField("col1", StringType(), nullable=True),  # 类型需与Cassandra表列完全匹配
    StructField("col2", StringType(), nullable=True),
    StructField("col3", StringType(), nullable=True),
    StructField("col4", StringType(), nullable=True),
    StructField("col5", IntegerType(), nullable=True),
    StructField("col6", StringType(), nullable=True),
    StructField("col7", StringType(), nullable=True)
])

# 将binary类型的value转为字符串后解析成结构化DataFrame
parsed_df = df.select(
    from_json(col("value").cast("string"), kafka_value_schema).alias("data")
).select("data.*")

parsed_df.printSchema()

# 写入Cassandra
ds = parsed_df.writeStream \
    .trigger(processingTime='15 seconds') \
    .format("org.apache.spark.sql.cassandra") \
    .option("checkpointLocation","{checkpoint}") \
    .options(table='{table}',keyspace="{key}") \
    .outputMode('update') \
    .start()

ds.awaitTermination()

报错信息

com.datastax.spark.connector.datasource.CassandraCatalogException: Attempting to write to C* Table but missing primary key columns: [col1,col2,col3]

at com.datastax.spark.connector.datasource.CassandraWriteBuilder.(CassandraWriteBuilder.scala:44)
at com.datastax.spark.connector.datasource.CassandraTable.newWriteBuilder(CassandraTable.scala:69)
at org.apache.spark.sql.execution.streaming.StreamExecution.createStreamingWrite(StreamExecution.scala:590)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.logicalPlan$lzycompute(MicroBatchExecution.scala:140)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.logicalPlan(MicroBatchExecution.scala:59)
at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:295)
at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStr
at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:209)

Traceback (most recent call last):

File "/home/test.py", line 33, in
ds.awaitTermination()

File "/venv/lib64/python3.6/site-packages/pyspark/python/lib/pyspark.zip/pyspark/sql/streaming.py", line 101, in awaitTe

File "/venv/lib64/python3.6/site-packages/pyspark/python/lib/py4j-0.10.9.5-src.zip/py4j/java_gateway.py", line 1322, in

File "/home/jeju/venv/lib64/python3.6/site-packages/pyspark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 117, in deco
pyspark.sql.utils.StreamingQueryException: Attempting to write to C* Table but missing primary key columns: [col1,col2,col3]
=== Streaming Query ===
Identifier: [id = d7da05f9-29a2-4597-a2c9-86a4ebfa65f2, runId = eea59c10-30fa-4939-8a30-03bd7c96b3f2]
Current Committed Offsets: {}
Current Available Offsets: {}

问题原因与解决方法

原因

Spark从Kafka读取的原始DataFrame仅包含Kafka自带的元数据字段(key、value、topic等),完全不包含Cassandra表所需的主键列col1、col2、col3及其他业务字段,直接写入触发Cassandra的主键校验规则,导致报错。

解决步骤

  1. 解析Kafka消息的value字段
    Kafka的业务数据通常存储在value字段中(二进制类型),需要将其转换为字符串后,根据消息格式(如JSON、Avro)解析成与Cassandra表结构匹配的结构化DataFrame。
    • 若使用Avro格式,需引入spark-avro依赖包进行解析。
    • 必须保证解析后的DataFrame字段名、数据类型与Cassandra表完全对应,尤其是主键列。
  2. 修正运行命令的配置项拼写
    原命令中com.datastax.spark:spark.cassandra.connectiohost为拼写错误,正确配置项为spark.cassandra.connection.host,否则会导致Cassandra连接失败。
  3. 使用解析后的DataFrame写入Cassandra
    将解析完成的结构化DataFrame传入写入流,确保包含Cassandra表的所有主键列和需要写入的业务字段。

额外注意事项

  • 确认Cassandra表的主键定义,若为复合主键则必须包含所有主键列。
  • 严格匹配Spark Schema与Cassandra表的字段类型,如Cassandra的text对应Spark的StringType,int对应IntegerType。

内容的提问来源于stack exchange,提问作者hi-inbeom

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:10:22