Spark容器无法连接/写入Cassandra容器表的问题求助
问题:Spark集群无法连接Docker网络中的Cassandra服务
我将Kafka Broker、Spark Master/Worker、Cassandra部署在同一Docker confluent 网络下,配置如下:
broker: image: confluentinc/cp-server:7.4.0 hostname: broker container_name: broker depends_on: zookeeper: condition: service_healthy ports: - "9092:9092" - "9101:9101" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1 KAFKA_CONFLUENT_BALANCER_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_JMX_PORT: 9101 KAFKA_JMX_HOSTNAME: localhost KAFKA_CONFLUENT_SCHEMA_REGISTRY_URL: http://schema-registry:8081 CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092 CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1 CONFLUENT_METRICS_ENABLE: 'false' CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous' networks: - confluent healthcheck: test: [ "CMD", "bash", "-c", 'nc -z localhost 9092' ] interval: 10s timeout: 5s retries: 5 spark-master: image: bitnami/spark:latest volumes: - ./spark_stream.py:/opt/bitnami/spark/spark_stream.py command: bin/spark-class org.apache.spark.deploy.master.Master ports: - "9090:8080" - "7077:7077" networks: - confluent spark-worker: image: bitnami/spark:latest command: bin/spark-class org.apache.spark.deploy.worker.Worker spark://spark-master:7077 depends_on: - spark-master environment: SPARK_MODE: worker SPARK_WORKER_CORES: 2 SPARK_WORKER_MEMORY: 1g SPARK_MASTER_URL: spark://spark-master:7077 networks: - confluent volumes: - ./spark_stream.py:/opt/bitnami/spark/spark_stream.py cassandra_db: image: cassandra:latest container_name: cassandra hostname: cassandra ports: - "9042:9042" environment: - MAX_HEAP_SIZE=512M - HEAP_NEWSIZE=100M - CASSANDRA_USERNAME=cassandra - CASSANDRA_PASSWORD=cassandra networks: - confluent # Define the networks networks: confluent:
代码与错误情况
- SparkSession配置:
s_conn = SparkSession.builder \ .appName('SparkDataStreaming') \ .config('spark.cassandra.connection.host', 'cassandra') \ .config("spark.cassandra.connection.port", "9042") \ .config("spark.cassandra.auth.username", "cassandra") \ .config("spark.cassandra.auth.password", "cassandra") \ .getOrCreate()
- 本地测试Cassandra连接(成功):
auth_provider = PlainTextAuthProvider(username='cassandra', password='cassandra') cluster = Cluster(['localhost'], port=9042, auth_provider=auth_provider, connect_timeout=60) cas_session = cluster.connect()
- Spark读取Kafka流:
spark_df = spark_conn.readStream \ .format('kafka') \ .option('kafka.bootstrap.servers', 'broker:29092') \ .option('subscribe', 'users_created') \ .option('startingOffsets', 'earliest') \ .load()
- 提交命令:
spark-submit --master spark://localhost:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=cassandra --conf spark.cassandra.connection.port=9042 spark_stream.py
错误1:使用cassandra作为主机名时
ERROR CassandraConnectorConf: Unknown host 'cassandra' java.net.UnknownHostException: cassandra: nodename nor servname provided, or not known
错误2:改用localhost/127.0.0.1时
ERROR TaskSetManager: Task 0 in stage 0.0 failed 4 times; aborting job 24/11/14 01:59:40 ERROR WriteToDataSourceV2Exec: Data source write support MicroBatchWrite[epoch: 0, writer: CassandraBulkWrite(...)] is aborting.
已确认Kafka主题有数据,排除列匹配问题,寻求解决方法。
解决方案
核心原因分析
- Spark Driver与Worker的网络隔离:本地执行
spark-submit时,Driver运行在宿主机,而Spark Worker在Docker容器的confluent网络内。- 用
cassandra主机名时,宿主机无法解析Docker网络内的域名,导致Driver找不到Cassandra。 - 用
localhost时,Worker容器内的localhost指向自身,而非宿主机映射的Cassandra服务,Worker无法连接。
- 用
- Cassandra认证配置缺失:SparkSession中的用户名密码配置可能未传递到Worker节点。
具体修复步骤
1. 统一网络访问方式
方案A:让Driver加入Docker网络
使用同版本Spark镜像在Docker内部提交任务,确保Driver和Worker在同一网络:
docker run --rm --network confluent -v $(pwd):/app bitnami/spark:latest spark-submit --master spark://spark-master:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=cassandra --conf spark.cassandra.connection.port=9042 --conf spark.cassandra.auth.username=cassandra --conf spark.cassandra.auth.password=cassandra /app/spark_stream.py
方案B:用宿主机Docker网络IP配置连接
- 执行
docker network inspect confluent获取宿主机在该网络中的IPv4地址(如172.18.0.1)。 - 修改SparkSession配置:
s_conn = SparkSession.builder \ .appName('SparkDataStreaming') \ .config('spark.cassandra.connection.host', '172.18.0.1') \ .config("spark.cassandra.connection.port", "9042") \ .config("spark.cassandra.auth.username", "cassandra") \ .config("spark.cassandra.auth.password", "cassandra") \ .getOrCreate() - 更新
spark-submit命令:spark-submit --master spark://localhost:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=172.18.0.1 --conf spark.cassandra.connection.port=9042 --conf spark.cassandra.auth.username=cassandra --conf spark.cassandra.auth.password=cassandra spark_stream.py
2. 确保认证配置全局生效
在spark-submit命令中显式添加认证参数,保证Worker节点能获取到用户名密码配置。
3. 验证网络连通性
在Spark Worker容器内测试Cassandra访问:
docker exec -it <spark-worker-container-id> ping cassandra docker exec -it <spark-worker-container-id> nc -zv cassandra 9042
若能ping通且端口开放,说明网络无问题,需检查Spark配置或提交方式。
内容的提问来源于stack exchange,提问作者user28291353
相关产品推荐
相关产品推荐

