如何使用Spark DataFrame列值更新对应Elasticsearch索引的文档
问题原因
你传入save()方法的"{esIndex}"会被Elasticsearch Spark连接器识别为固定的索引名字面量,不会和DataFrame的esIndex字段做绑定,因此无法动态读取每行的索引取值,才会提示{esIndex}索引不存在。
解决方案
新增es.mapping.index配置项,指定用DataFrame中的esIndex字段作为每行数据写入的目标索引即可,修改后的命令4代码如下:
import org.apache.spark.sql.functions.col val retval = df.toDF.write .format("org.elasticsearch.spark.sql") .option("es.nodes.wan.only","true") .option("es.mapping.id","id") .option("es.mapping.index", "esIndex") // 指定用esIndex字段作为每行的目标写入索引 .option("es.mapping.exclude", "id, esIndex") .option("es.port","80") .option("es.nodes", esUrl) .option("es.write.operation", "update") .option("es.index.auto.create", "no") .mode("Append") .save()
修改后连接器会自动读取每行esIndex字段的取值作为写入目标索引,且你原本配置的es.mapping.exclude已经排除了esIndex字段,该字段不会被写入到ES文档内容中,符合需求。
内容的提问来源于stack exchange,提问作者Mohammed-AR
相关产品推荐
相关产品推荐

