Docker-Compose部署Kafka+Spark时PySpark读取Kafka数据源失败问题
我需要将Windows本地的TXT文件发送至Docker容器中的Kafka,再由另一容器内的PySpark消费并做map()转换处理。通过docker-compose定义了自定义网络,部署了spark-master、两个spark-worker、ZooKeeper和Kafka容器,因Kafka、Spark与Python版本兼容问题选用了最新的Bitnami镜像。
已创建名为demo的主题,通过Kafka容器执行命令bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic demo发送文本,但在spark-master容器中执行spark-submit mycode.py时出现错误:
Traceback (most recent call last): File "/src/structuredKafkaSpark.py", line 12, in <module> df = spark \ File "/opt/bitnami/spark/python/lib/pyspark.zip/pyspark/sql/streaming.py", line 469, in load File "/opt/bitnami/spark/python/lib/py4j-0.10.9.5-src.zip/py4j/java_gateway.py", line 1321, in __call__ File "/opt/bitnami/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 196, in deco pyspark.sql.utils.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".
遵循Spark官方指南操作后仍未解决问题,以下是我的docker-compose配置文件和PySpark代码:
docker-compose配置文件
version: "3.7" networks: datapipeline: driver: bridge services: spark-master: build: context: ./spark dockerfile: ./Dockerfile container_name: "spark-master" environment: - SPARK_MODE=master - SPARK_LOCAL_IP=spark-master - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no ports: - "7077:7077" - "8080:8080" volumes: - ./src:/src - ./data:/data - ./output:/output networks: - datapipeline spark-worker-1: image: docker.io/bitnami/spark:latest container_name: "spark-worker-1" environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077 - SPARK_WORKER_MEMORY=2G - 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-worker-2: image: docker.io/bitnami/spark:latest container_name: "spark-worker-2" environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077 - SPARK_WORKER_MEMORY=2G - SPARK_WORKER_CORES=1 - SPARK_RPC_AUTHENTICATION_ENABLED=no - SPARK_RPC_ENCRYPTION_ENABLED=no - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no - SPARK_SSL_ENABLED=no # ----------------- # # Apache Kafka # # ----------------- # zookeeper: image: docker.io/bitnami/zookeeper:latest container_name: "zookeeper" ports: - "2181:2181" environment: - ALLOW_ANONYMOUS_LOGIN=yes networks: - datapipeline kafka: image: docker.io/bitnami/kafka:latest container_name: "kafka" ports: - "9092:9092" environment: - KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181 - ALLOW_PLAINTEXT_LISTENER=yes depends_on: - zookeeper volumes: - ./producer:/producer networks: - datapipeline
PySpark代码
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, split spark = SparkSession \ .builder \ .appName("StructuredNetworkWordCount") \ .config("spark.driver.host", "localhost")\ .getOrCreate() # Create DataFrame representing the stream of input lines from kafka df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "kafka:9092") \ .option("subscribe", "demo") \ .load() df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") # Split the lines into words words = df.select( explode( split(df.value, " ") ).alias("word") ) # Generate running word count wordCounts = words.groupBy("word").count() # Start running the query that prints the running counts to the console query = wordCounts \ .writeStream \ .outputMode("update") \ .format("console") \ .start() query.awaitTermination()
核心问题:缺失Spark-Kafka连接器
Bitnami的Spark镜像默认未包含Kafka连接器依赖,导致Spark无法识别kafka数据源,提供两种解决方式:
方式1:提交作业时指定依赖包
执行spark-submit时,通过--packages参数拉取与Spark版本匹配的连接器包。先通过spark-submit --version查看容器内Spark版本,比如Spark 3.5.x对应的命令:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0 /src/structuredKafkaSpark.py
注:
_2.12为Scala版本,Bitnami Spark镜像默认使用的Scala版本可通过容器内Spark配置确认。
方式2:构建自定义镜像预安装连接器
在./spark/Dockerfile中添加下载连接器的步骤,示例如下:
FROM docker.io/bitnami/spark:latest # 下载对应Spark版本的连接器包 RUN wget -P /opt/bitnami/spark/jars/ https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_2.12/3.5.0/spark-sql-kafka-0-10_2.12-3.5.0.jar RUN wget -P /opt/bitnami/spark/jars/ https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/3.4.0/kafka-clients-3.4.0.jar
重新构建镜像:docker-compose build spark-master,再启动容器即可。
额外优化点
- Kafka生产者命令优化:在Kafka容器内执行生产者命令时,使用容器名作为bootstrap地址更标准:
bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic demo
- 本地TXT文件发送到Kafka:可在Windows本地运行Python脚本,直接连接
localhost:9092(Kafka容器已映射端口到主机)发送文件内容:
from kafka import KafkaProducer import time producer = KafkaProducer(bootstrap_servers='localhost:9092') with open('本地文件路径.txt', 'r', encoding='utf-8') as f: for line in f: producer.send('demo', value=line.strip().encode('utf-8')) time.sleep(0.1) producer.flush()
- Spark代码修正:原代码中
df.selectExpr的结果未赋值,后续操作仍使用原始二进制类型的value会报错,修改如下:
df_str = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") # Split the lines into words words = df_str.select( explode( split(df_str.value, " ") ).alias("word") )
内容的提问来源于stack exchange,提问作者yaviens

