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

Spark+Kafka应用写入Cassandra时遭遇ClassNotFoundException求助

问题:Spark Streaming写入Cassandra触发ClassNotFoundException异常

环境配置

  • Python 3.6.8
  • Spark 3.1.2
  • Kafka 2.8.1
  • Cassandra 3.0.8
  • 代码部署IP与Spark Master节点不同

已成功通过Spark Streaming读取Kafka的test0003主题,但调用writeStream将数据写入Cassandra的testkey.testtable时失败,触发ClassNotFoundException: Failed to find data source: org.apache.spark.sql.cassandra异常。

运行命令

spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 test.py

相关代码

from pyspark.sql import SparkSession

while True:
    # Spark Bridge local to master == Connect master
    spark = SparkSession.builder \
        .master("spark://masterIp:7077") \
        .appName("Spark_Streaming+kafka+cassandra") \
        .getOrCreate()

    # Read Stream From test0003 at BootStrap
    df = spark.readStream \
      .format("kafka") \
      .option("kafka.bootstrap.servers", "DBip:9092") \
      .option('startingOffsets','earliest') \
      .option('failOnDataLoss','False') \
      .option("subscribe", "test0003") \
      .load()

    df.printSchema()

    # write Stream at cassandra
    ds = df \
      .writeStream \
      .trigger(processingTime='15 seconds') \
      .format("org.apache.spark.sql.cassandra") \
      .outputMode('append') \
      .option("checkpointLocation","./kafka") \
      .option("kafka.bootstrap.servers", "master_Ip:9092") \
      .option("keyspace","testkey") \
      .option("table","testtable") \
      .start()

    break

报错信息

Traceback (most recent call last):
  File "./test.py", line 30, in <module>
    .option("table","testtable") \
  File "./venv/lib64/python3.6/site-packages/pyspark/python/lib/pyspark.zip/pyspark/sql/streaming.py", line 1202, in start
  File "./venv/lib64/python3.6/site-packages/pyspark/python/lib/py4j-0.10.9.5-src.zip/py4j/java_gateway.py", line 1322, in __call__
  File "./venv/lib64/python3.6/site-packages/pyspark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 111, in deco
  File "./venv/lib64/python3.6/site-packages/pyspark/python/lib/py4j-0.10.9.5-src.zip/py4j/protocol.py", line 328, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o47.start.
: java.lang.ClassNotFoundException:
Failed to find data source: org.apache.spark.sql.cassandra. Please find packages at

已找到相关解决思路,但问题仍未解决。本人是入职1个月的初级开发者,首次求助,希望得到帮助。


内容的提问来源于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 07:35:18