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

Spark 2.4.0 Structured Streaming写入Cassandra无数据问题求助

排查Spark 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 07:50:35