Docker容器中Spark与Cassandra连接失败问题求助
问题排查:Spark-Cassandra连接器数据源找不到错误
问题描述
通过docker-compose部署Spark(master/worker)与Cassandra容器,Cassandra已导入数据。将spark-cassandra-connector_2.12-latest_version.jar放入Spark容器jars目录后,运行PySpark读取脚本时出现[DATA_SOURCE_NOT_FOUND]错误,提示找不到org.apache.spark.sql.cassandra数据源。Spark Hello World程序可正常运行,所有容器处于同一Docker网络。
环境配置
docker-compose配置
spark: image: docker.io/bitnami/spark:3.4 container_name: spark_master environment: - SPARK_MODE=master - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no - SPARK_USER=spark ports: - '8081:8081' spark-worker: image: docker.io/bitnami/spark:3.4 container_name: spark_worker environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark:7077 - SPARK_WORKER_MEMORY=1G - SPARK_WORKER_CORES=1 - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no - SPARK_USER=spark cassandra: container_name: cassandra-container image: cassandra:latest ports: - "7000:7000" - "7001:7001" - "7199:7199" - "9042:9042" - "9160:9160" volumes: - cassandra-data:/var/lib/cassandra restart: always
PySpark连接脚本
from pyspark.sql import SparkSession # Define your Cassandra connection settings cassandra_host = "127.0.0.1" # Replace with the actual IP address of your Cassandra container cassandra_port = "9042" # Default Cassandra port cassandra_keyspace = "mykey_space" # Replace with your Cassandra keyspace cassandra_table = "test_table" # Replace with your Cassandra table # Create a Spark session spark = SparkSession.builder \ .appName("CassandraConnectionTest") \ .config("spark.cassandra.connection.host", cassandra_host) \ .config("spark.cassandra.connection.port", cassandra_port) \ .getOrCreate() # Test the Cassandra connection try: df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .options(table=cassandra_table,keyspace= cassandra_keyspace) \ .load() df.show(5) print("Cassandra connection test successful.") except Exception as e: print(f"Error: {str(e)}") print("Cassandra connection test failed.") # Stop the Spark session spark.stop()
运行命令
docker exec -it spark_worker /bin/bash -c "spark-submit --jars /opt/bitnami/spark/jars/spark-cassandra-connector_2.12-latest_version.jar /opt/bitnami/spark/spark-script/cassandra_spark.py"
Docker网络信息
[{ "Name": "project_default", "Id": "c28118c6ba0511ec9f54a97a6d7b17113a75995d4bfed6d18fc80cf3616acb03", "Created": "2023-09-18T12:19:08.729378699Z", "Scope": "local", "Driver": "bridge", "EnableIPv6": false, "IPAM": { "Driver": "default", "Options": null, "Config": [{ "Subnet": "172.20.0.0/16", "Gateway": "172.20.0.1" }] }, "Internal": false, "Attachable": false, "Ingress": false, "ConfigFrom": { "Network": "" }, "ConfigOnly": false, "Containers": { "89396ef426b166e4d293011ef2c8c8ff60007d2baeba2734ef2e91a0cf7f60ec": { "Name": "cassandra-container", "EndpointID": "5a3804e47d659bcdf19a44477e85fc7c250b52b0375d342984c7d0033ff1b6cc", "MacAddress": "02:42:ac:14:00:04", "IPv4Address": "172.20.0.4/16", "IPv6Address": "" }, "95cdbc38d6717b0318e81e33dfe990dc5da461d22f60976d132559d568c9365f": { "Name": "spark_worker", "EndpointID": "cc6b69c4c2f2de02965efc78a53dc92d2f06e4e2a0d1e9dd29fb161308032a26", "MacAddress": "02:42:ac:14:00:03", "IPv4Address": "172.20.0.3/16", "IPv6Address": "" }, "e4864bdc0ff72f0185aceaf6f4844ef75ec45f6673cc776529eb8ca2f67943c6": { "Name": "spark_master", "EndpointID": "83c4b290a2730cec39849189b326513b33b70042832dacf9c5b47b2d92f90351", "MacAddress": "02:42:ac:14:00:02", "IPv4Address": "172.20.0.2/16", "IPv6Address": "" } }, "Options": {}, "Labels": { "com.docker.compose.network": "default", "com.docker.compose.project": "project", "com.docker.compose.version": "2.18.1" } }]
错误日志
23/09/17 17:45:10 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable 23/09/17 17:45:11 INFO SparkContext: Running Spark version 3.4.1 ... Error: An error occurred while calling o34.load. : org.apache.spark.SparkClassNotFoundException: [DATA_SOURCE_NOT_FOUND] Failed to find the data source: org.apache.spark.sql.cassandra. Please find packages at `https://spark.apache.org/third-party-projects.html`. ... Caused by: java.lang.ClassNotFoundException: org.apache.spark.sql.cassandra.DefaultSource ...
排查与解决方案
1. 确认连接器版本兼容性
Spark 3.4.1对应Scala 2.12,必须使用匹配版本的spark-cassandra-connector,正确版本为3.4.0,文件名应为spark-cassandra-connector_2.12-3.4.0.jar。避免使用latest_version这类模糊命名,防止下载到不兼容版本。
2. 检查jar包路径与权限
- 执行以下命令确认jar包存在于指定目录:
docker exec spark_worker ls /opt/bitnami/spark/jars/ | grep cassandra - 若jar包存在但权限异常,修正权限:
docker exec spark_worker chown spark:spark /opt/bitnami/spark/jars/spark-cassandra-connector_2.12-3.4.0.jar - 或者直接省略
--jars参数,将jar包放在Spark默认jars目录,Spark会自动加载。
3. 修正Cassandra连接地址
容器内部无法通过127.0.0.1访问Cassandra,需替换为Docker网络内的容器名cassandra-container或内网IP172.20.0.4,修改PySpark脚本中的cassandra_host:
cassandra_host = "cassandra-container"
4. 完善SparkSession配置
在构建SparkSession时添加连接器包配置,让Spark自动管理依赖(需容器联网):
spark = SparkSession.builder \ .appName("CassandraConnectionTest") \ .config("spark.cassandra.connection.host", "cassandra-container") \ .config("spark.jars.packages", "com.datastax.spark:spark-cassandra-connector_2.12:3.4.0") \ .getOrCreate()
此时运行命令无需指定--jars参数:
docker exec -it spark_worker /bin/bash -c "spark-submit /opt/bitnami/spark/spark-script/cassandra_spark.py"
内容的提问来源于stack exchange,提问作者Amine SAIHI
相关产品推荐
相关产品推荐

