Spark 2.4.0 Structured Streaming写入Cassandra无数据问题求助
嘿,我帮你排查下这个Spark流写Cassandra没数据的问题,结合你的Spark 2.4.0版本和代码,大概率是这几个点出了问题:
1. 字段类型不匹配(最核心原因)
你虽然定义了schema,但实际解析Kafka JSON数据时根本没用到它——用get_json_object提取出来的所有字段都是字符串类型,但Cassandra是强类型数据库,如果你的sensor表中humidity是int、time是timestamp这类类型,Spark写入时会因为类型不匹配静默失败(不会抛出明显错误,但数据直接被丢弃)。
两种修正方式:
方式一:用预定义的schema解析JSON
直接用你写好的schema来解析Kafka的value字符串,自动匹配类型:
df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafkanode") \ .option("subscribe", "iot-data-sensor") \ .load() \ .select(col("value").cast("string").alias("json_str")) \ .from_json("json_str", schema)
方式二:手动转换字段类型
如果坚持用get_json_object,要把每个字段转成对应类型:
from pyspark.sql.types import IntegerType, TimestampType df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafkanode") \ .option("subscribe", "iot-data-sensor") \ .load() \ .select( get_json_object(col("value").cast("string"), "$.humidity").cast(IntegerType()).alias("humidity"), get_json_object(col("value").cast("string"), "$.time").cast(TimestampType()).alias("time"), get_json_object(col("value").cast("string"), "$.temperature").cast(IntegerType()).alias("temperature"), get_json_object(col("value").cast("string"), "$.ph").cast(IntegerType()).alias("ph"), get_json_object(col("value").cast("string"), "$.sensor").alias("sensor"), get_json_object(col("value").cast("string"), "$.id").alias("id") )
2. 程序启动后直接退出,没来得及处理数据
你的代码最后只调用了.start(),但没加awaitTermination()——这会导致Spark流程序刚启动就立刻退出,根本没机会处理Kafka的数据并写入Cassandra。
修正很简单,在启动流之后加上等待:
stream = df.writeStream \ .foreachBatch(writeToCassandra) \ .outputMode("update") \ .start() stream.awaitTermination() # 让程序持续运行,处理流数据
3. Cassandra连接地址的空格坑
你在writeToCassandra里写的spark.cassandra.connection.host是"cassnode1, cassnode2"——逗号后面多了个空格!这会让Spark解析节点地址出错,连不上Cassandra集群,自然写不进数据。
改成这样:
.options("spark.cassandra.connection.host", "cassnode1,cassnode2")
4. Output Mode可能不符合你的场景
你用的是outputMode("update"),这个模式只有当数据有更新(比如某个key对应的字段值变化)时才会触发foreachBatch。如果你的Kafka数据全是新的、无重复的记录,可能不会触发写入逻辑。
如果你的需求是写入所有新数据,建议改成outputMode("append"):
df.writeStream \ .foreachBatch(writeToCassandra) \ .outputMode("append") \ .start()
5. 检查Spark-Cassandra Connector版本兼容性
Spark 2.4.0必须搭配对应版本的Connector,建议用2.4.0版本的connector(比如Maven依赖:com.datastax.spark:spark-cassandra-connector_2.11:2.4.0),版本不兼容也会导致写入失败。
额外调试小技巧
如果还是找不到问题,可以在writeToCassandra里加个打印,确认是不是根本没进入写入逻辑:
def writeToCassandra(writeDF, epochId): print(f"正在处理批次 {epochId},记录数:{writeDF.count()}") # 打印批次信息 writeDF.write \ .format("org.apache.spark.sql.cassandra") \ .mode('append') \ .options("spark.cassandra.connection.host", "cassnode1,cassnode2") \ .options(table="sensor", keyspace="sensordb") \ .save()
如果控制台看不到这个打印,说明foreachBatch没被触发,就得回头查Output Mode或者Kafka数据的问题;如果能看到打印但Cassandra还是没数据,那就是写入Cassandra的环节出了问题(比如权限、表结构不匹配)。
内容的提问来源于stack exchange,提问作者Fidia Rosianti

