在Jupyter Notebook中用PySpark从Kafka流数据到Elasticsearch失败求助
问题:PySpark流式写入Elasticsearch失败(写入HDFS正常)
在Jupyter Notebook中通过PySpark从Kafka Topic流式读取数据,写入HDFS时一切正常,但尝试写入Elasticsearch索引时执行失败,抛出NoSuchMethodError异常。
报错信息
Driver stacktrace:) 24/09/14 04:17:08 ERROR Executor: Exception in task 39.0 in stage 15.0 (TID 193) java.lang.NoSuchMethodError: 'org.apache.spark.sql.catalyst.encoders.ExpressionEncoder org.apache.spark.sql.catalyst.encoders.RowEncoder$.apply(org.apache.spark.sql.types.StructType)' at org.elasticsearch.spark.sql.streaming.EsStreamQueryWriter.<init>(EsStreamQueryWriter.scala:50) at org.apache.spark.sql.streaming.EsSparkSqlStreamingSink.$anonfun$addBatch$5(EsSparkSqlStreamingSink.scala:72) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833)
版本信息
- Spark Version: 3.5.0
- Elasticsearch: docker.elastic.co/elasticsearch/elasticsearch:8.9.0
Elasticsearch及Kibana Docker配置
elasticsearch: image: docker.elastic.co/elasticsearch/elasticsearch:8.9.0 container_name: elasticsearch environment: - discovery.type=single-node - ELASTIC_PASSWORD=password#123 - xpack.security.enabled=false - xpack.security.transport.ssl.enabled=false ports: - '9200:9200' - '9300:9300' kibana: image: docker.elastic.co/kibana/kibana:8.9.0 container_name: kibana ports: - '5601:5601' environment: - ELASTICSEARCH_HOSTS=http://elasticsearch:9200 - xpack.security.enabled=false depends_on: - elasticsearch
故障代码(PySpark)
# Create Spark Session spark = SparkSession.builder \ .appName("KafkaToHDFS") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0,org.elasticsearch:elasticsearch-spark-30_2.12:8.9.0") \ .config("es.nodes", "192.XXX.XX.144") \ .config("es.port", "9200") \ .config("es.nodes.wan.only", "true") \ .config("spark.sql.adaptive.enabled", "false") \ .getOrCreate() es_query = person_df.writeStream \ .format("es") \ .queryName("writing_to_es") \ .option("es.nodes", "192.XXX.XX.144:9200") \ .option("es.resource", "uc_person_plot/_doc") \ .option("checkpointLocation", "hdfs://namenode:9000/uc/es/checkpoint_dir") \ .outputMode("append") \ .start()
解决方案
1. 对齐依赖包与Spark版本
报错核心原因是Spark 3.5.0的API变更,导致旧版本的Elasticsearch Spark Connector和Kafka Connector无法兼容:
- Kafka Connector版本必须与Spark版本完全一致:将
spark-sql-kafka-0-10_2.12:3.3.0替换为spark-sql-kafka-0-10_2.12:3.5.0 - Elasticsearch Spark Connector 8.9.0不支持Spark 3.5,需升级到8.11.3及以上版本(该版本开始官方支持Spark 3.5),将
elasticsearch-spark-30_2.12:8.9.0替换为elasticsearch-spark-30_2.12:8.11.3
修改后的SparkSession配置:
spark = SparkSession.builder \ .appName("KafkaToHDFS") \ .config("spark.jars.packages", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0,org.elasticsearch:elasticsearch-spark-30_2.12:8.11.3") \ .config("es.nodes", "192.XXX.XX.144") \ .config("es.port", "9200") \ .config("es.nodes.wan.only", "true") \ .config("spark.sql.adaptive.enabled", "false") \ .getOrCreate()
2. 简化Elasticsearch写入配置
SparkSession已全局配置es.nodes和es.port,写入流时无需重复设置,可简化代码:
es_query = person_df.writeStream \ .format("es") \ .queryName("writing_to_es") \ .option("es.resource", "uc_person_plot/_doc") \ .option("checkpointLocation", "hdfs://namenode:9000/uc/es/checkpoint_dir") \ .outputMode("append") \ .start()
3. 验证Elasticsearch连接
确保192.XXX.XX.144:9200可正常访问,容器端口映射正确,且ES的xpack.security.enabled已关闭(当前配置已满足)。
内容的提问来源于stack exchange,提问作者Afzal Abdul Azeez
相关产品推荐
相关产品推荐

