Spark+Kafka流写入Cassandra报错:缺失主键列[col1,col2,col3]
运行环境
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的主键校验规则,导致报错。
解决步骤
- 解析Kafka消息的value字段
Kafka的业务数据通常存储在value字段中(二进制类型),需要将其转换为字符串后,根据消息格式(如JSON、Avro)解析成与Cassandra表结构匹配的结构化DataFrame。- 若使用Avro格式,需引入
spark-avro依赖包进行解析。 - 必须保证解析后的DataFrame字段名、数据类型与Cassandra表完全对应,尤其是主键列。
- 若使用Avro格式,需引入
- 修正运行命令的配置项拼写
原命令中com.datastax.spark:spark.cassandra.connectiohost为拼写错误,正确配置项为spark.cassandra.connection.host,否则会导致Cassandra连接失败。 - 使用解析后的DataFrame写入Cassandra
将解析完成的结构化DataFrame传入写入流,确保包含Cassandra表的所有主键列和需要写入的业务字段。
额外注意事项
- 确认Cassandra表的主键定义,若为复合主键则必须包含所有主键列。
- 严格匹配Spark Schema与Cassandra表的字段类型,如Cassandra的
text对应Spark的StringType,int对应IntegerType。
内容的提问来源于stack exchange,提问作者hi-inbeom

