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
相关产品推荐
相关产品推荐

