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

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,再启动容器即可。

额外优化点

  1. Kafka生产者命令优化:在Kafka容器内执行生产者命令时,使用容器名作为bootstrap地址更标准:
bin/kafka-console-producer.sh --bootstrap-server kafka:9092 --topic demo
  1. 本地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()
  1. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 00:15:25