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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 14:14:58