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
相关产品推荐
相关产品推荐

