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

Docker环境下PySpark连接Kafka报错:Failed to find data source: kafka

问题:Spark连接Kafka时提示找不到数据源

环境配置(docker-compose.yml)

version: '2'
services:
  zookeeper:
    image: quay.io/debezium/zookeeper:${DEBEZIUM_VERSION}
    ports:
     - 2181:2181
     - 2888:2888
     - 3888:3888
  kafka:
    image: quay.io/debezium/kafka:${DEBEZIUM_VERSION}
    ports:
     - 9092:9092
    links:
     - zookeeper
    environment:
     - ZOOKEEPER_CONNECT=zookeeper:2181
  mysql:
    image: quay.io/debezium/example-mysql:${DEBEZIUM_VERSION}
    ports:
     - 3306:3306
    environment:
     - MYSQL_ROOT_PASSWORD=debezium
     - MYSQL_USER=mysqluser
     - MYSQL_PASSWORD=mysqlpw
  connect:
    image: quay.io/debezium/connect:${DEBEZIUM_VERSION}
    ports:
     - 8083:8083
    links:
     - kafka
     - mysql
    environment:
     - BOOTSTRAP_SERVERS=kafka:9092
     - GROUP_ID=1
     - CONFIG_STORAGE_TOPIC=my_connect_configs
     - OFFSET_STORAGE_TOPIC=my_connect_offsets
     - STATUS_STORAGE_TOPIC=my_connect_statuses
  spark-master:
      image: docker.io/bitnami/spark:3.3
      environment:
        - SPARK_MODE=master
        - SPARK_RPC_AUTHENTICATION_ENABLED=no
        - SPARK_RPC_ENCRYPTION_ENABLED=no
        - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
        - SPARK_SSL_ENABLED=no
      ports:
        - '8080:8080'
  spark-worker:
    image: docker.io/bitnami/spark:3.3
    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
    links:
      - kafka
  jupyter:
    image: jupyter/pyspark-notebook
    environment:
      - GRANT_SUDO=yes
      - JUPYTER_ENABLE_LAB=yes
      - JUPYTER_TOKEN=mysecret
    ports:
      - "8888:8888"
    volumes:
      - /Users/eugenegoldberg/jupyter_notebooks:/home/eugene
    depends_on:
      - spark-master  

PySpark连接Kafka代码

from pyspark import SparkConf, SparkContext
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *

import os

spark_version = '3.3.1'
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-sql-kafka-0-10_2.12:{}'.format(spark_version)

packages = [
    f'org.apache.kafka:kafka-clients:3.3.1'
]

# Create SparkSession
spark = SparkSession.builder \
    .appName("Kafka Streaming Example") \
    .config("spark.driver.host", "host.docker.internal") \
    .config("spark.jars.packages", ",".join(packages)) \
    .getOrCreate()

# Define the Kafka topic and Kafka server/port
topic = "dbserver1.inventory.customers"
kafkaServer = "kafka:9092" # assuming kafka is running on a container named 'kafka'

# Read data from kafka topic
df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", kafkaServer) \
  .option("subscribe", topic) \
  .load()

报错信息

---------------------------------------------------------------------------
AnalysisException                         Traceback (most recent call last)
Cell In[9], line 32
     24 kafkaServer = "kafka:9092" # assuming kafka is running on a container named 'kafka'
     26 # Read data from kafka topic
     27 df = spark \
     28   .readStream \
     29   .format("kafka") \
     30   .option("kafka.bootstrap.servers", kafkaServer) \
     31   .option("subscribe", topic) \
---> 32   .load()

File /usr/local/spark/python/pyspark/sql/streaming.py:469, in DataStreamReader.load(self, path, format, schema, **options)
    467     return self._df(self._jreader.load(path))
    468 else:
--> 469     return self._df(self._jreader.load())

File /usr/local/spark/python/lib/py4j-0.10.9.5-src.zip/py4j/java_gateway.py:1321, in JavaMember.__call__(self, *args)
   1315 command = proto.CALL_COMMAND_NAME +\
   1316     self.command_header +\
   1317     args_command +\
   1318     proto.END_COMMAND_PART
   1320 answer = self.gateway_client.send_command(command)
-> 1321 return_value = get_return_value(
   1322     answer, self.gateway_client, self.target_id, self.name)
   1324 for temp_arg in temp_args:
   1325     temp_arg._detach()

File /usr/local/spark/python/pyspark/sql/utils.py:196, in capture_sql_exception.<locals>.deco(*a, **kw)
    192 converted = convert_exception(e.java_exception)
    193 if not isinstance(converted, UnknownException):
    194     # Hide where the exception came from that shows a non-Pythonic
    195     # JVM exception message.
--> 196     raise converted from None
    197 else:
    198     raise

AnalysisException:  Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".    

解决方案

1. 修正Spark依赖包配置

报错核心原因是缺少Spark连接Kafka的核心依赖spark-sql-kafka-0-10_2.12,同时代码中存在依赖配置重复的问题,修改如下:

from pyspark import SparkConf, SparkContext
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *

spark_version = '3.3.1'
# 包含Spark Kafka核心依赖和kafka-clients,版本与Spark匹配
packages = [
    f'org.apache.spark:spark-sql-kafka-0-10_2.12:{spark_version}',
    'org.apache.kafka:kafka-clients:3.3.1'
]

# Create SparkSession
spark = SparkSession.builder \
    .appName("Kafka Streaming Example") \
    .config("spark.driver.host", "host.docker.internal") \
    .config("spark.jars.packages", ",".join(packages)) \
    .getOrCreate()

# 后续代码不变
topic = "dbserver1.inventory.customers"
kafkaServer = "kafka:9092"

df = spark \
  .readStream \
  .format("kafka") \
  .option("kafka.bootstrap.servers", kafkaServer) \
  .option("subscribe", topic) \
  .load()

同时删除原代码中os.environ['PYSPARK_SUBMIT_ARGS']这一行,避免与spark.jars.packages的配置冲突。

2. 确保Jupyter容器能访问Kafka服务

在docker-compose.yml的jupyter服务中添加Kafka的链接配置,保证容器能解析kafka主机名:

jupyter:
    image: jupyter/pyspark-notebook
    environment:
      - GRANT_SUDO=yes
      - JUPYTER_ENABLE_LAB=yes
      - JUPYTER_TOKEN=mysecret
    ports:
      - "8888:8888"
    volumes:
      - /Users/eugenegoldberg/jupyter_notebooks:/home/eugene
    depends_on:
      - spark-master
      - kafka  # 添加依赖
    links:
      - kafka  # 添加链接,确保主机名解析

3. 验证版本兼容性

确保spark-sql-kafka-0-10_2.12的版本与Spark版本完全一致(这里为3.3.1),kafka-clients版本建议与Debezium使用的Kafka版本兼容,当前3.3.1版本与Spark 3.3.1适配性良好。

内容的提问来源于stack exchange,提问作者Eugene Goldberg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 01:40:24