PySpark DataFrame写入Kafka失败求助:NoSuchMethodError异常
PySpark写入Kafka报错排查与解决
问题描述
尝试将简单PySpark DataFrame写入Kafka,相关配置已完成但持续报错。普通Python生产者脚本可正常运行,PySpark读流功能也无异常,唯独PySpark写Kafka失败。
代码示例
from pyspark.sql import SparkSession import os os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-streaming-kafka-0-10_2.12:3.2.0,org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0 pyspark-shell' # 创建SparkSession spark = SparkSession.builder \ .appName("KafkaSinkExample") \ .getOrCreate() # 示例数据 data = [("key1", "value1"), ("key2", "value2"), ("key3", "value3")] df = spark.createDataFrame(data, ["key", "value"]) # 将数据写入Kafka (df.write \ .format("kafka") \ .option("kafka.bootstrap.servers", "broker:29092") \ .option("topic", "test") \ .save())
Docker Compose配置
services: zookeeper: image: confluentinc/cp-zookeeper:7.5.0 hostname: zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 broker: image: confluentinc/cp-server:7.5.0 hostname: broker container_name: broker depends_on: - zookeeper 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: 'true' CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous' pyspark: image: jupyter/pyspark-notebook:latest container_name: pyspark ports: - "8888:8888" environment: - PYSPARK_PYTHON=python3 - PYSPARK_DRIVER_PYTHON=python3 - SPARK_HOME=/usr/local/spark volumes: - ./notebooks:/home/jovyan/work # 挂载你的notebooks文件夹,按需调整
报错信息
23/12/25 16:50:56 ERROR TaskSetManager: Task 1 in stage 0.0 failed 1 times; aborting job Traceback (most recent call last): (0 + 1) / 2] File "/home/jovyan/preprocessing/bing.py", line 23, in <module> .save()) ^^^^^^ File "/usr/local/spark/python/pyspark/sql/readwriter.py", line 1461, in save self._jwrite.save() File "/usr/local/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in __call__ File "/usr/local/spark/python/pyspark/errors/exceptions/captured.py", line 179, in deco return f(*a, **kw) ^^^^^^^^^^^ File "/usr/local/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o49.save. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 1 in stage 0.0 failed 1 times, most recent failure: Lost task 1.0 in stage 0.0 (TID 1) (0f08baeba383 executor driver): java.lang.NoSuchMethodError: 'boolean org.apache.spark.sql.catalyst.expressions.Cast$.apply$default$4()'
解决方案
这个NoSuchMethodError是Spark版本与Kafka连接器版本不兼容导致的核心问题:你指定的Kafka连接器是3.2.0版本,但jupyter/pyspark-notebook:latest镜像默认安装的Spark版本大概率高于3.2.0,两者API不匹配引发错误。
具体修复步骤:
- 确认容器内Spark版本:进入pyspark容器执行
spark-submit --version,查看实际运行的Spark版本。 - 同步连接器版本:将
PYSPARK_SUBMIT_ARGS中的连接器版本改为与容器内Spark一致的版本。例如如果Spark是3.5.0,修改为:os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages org.apache.spark:spark-streaming-kafka-0-10_2.12:3.5.0,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0 pyspark-shell' - 数据类型校验:写入Kafka要求
key和value列必须是字符串或二进制类型,你的示例数据符合要求,后续若有其他类型数据需提前转换。 - 可选版本绑定:如果不想手动匹配版本,可直接使用指定Spark版本的镜像,比如
jupyter/pyspark-notebook:spark-3.2.0,这样连接器用3.2.0即可正常工作。
内容的提问来源于stack exchange,提问作者kreemo
相关产品推荐
相关产品推荐

