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

如何使用PySpark将MongoDB集合数据同步到Elasticsearch,求代码示例

PySpark 同步 MongoDB 数据到 Elasticsearch 完整实现方案

前置依赖说明

  • 提前确认运行环境可同时连通 MongoDB 实例和 Elasticsearch 实例
  • 确认 MongoDB 连接字符串持有对应集合的读权限,Elasticsearch 持有目标索引的写权限
  • 提交任务时通过--packages参数指定对应版本的MongoDB连接器、ES连接器即可,无需提前手动下载

完整代码示例

from pyspark.sql import SparkSession

if __name__ == "__main__":
    # 初始化SparkSession,统一配置连接器参数
    spark = SparkSession.builder \
        .appName("MongoToESSync") \
        # MongoDB配置:直接填入你持有的连接字符串,注意末尾指定到具体库+集合,权限参数可跟在URI后
        .config("spark.mongodb.input.uri", "mongodb://[账号]:[密码]@[IP]:[端口]/[库名].[集合名]?authSource=admin") \
        # ES基础配置
        .config("spark.es.nodes", "ES节点IP,多个用逗号分隔") \
        .config("spark.es.port", "ES端口,默认9200") \
        # 若ES开启身份认证,打开下方两个配置填入对应信息
        # .config("spark.es.net.http.auth.user", "ES账号")
        # .config("spark.es.net.http.auth.pass", "ES密码") \
        # 跨网段/云ES服务必须开启该配置,关闭节点自动发现
        .config("spark.es.nodes.wan.only", "true") \
        .getOrCreate()

    # 步骤1:从MongoDB读取数据
    df = spark.read.format("mongo").load()

    # 可选:数据处理逻辑
    # 字段重命名示例:把MongoDB的_id转为mongo_id避免和ES的_id冲突
    # df = df.withColumnRenamed("_id", "mongo_id")
    # 过滤示例:全量同步过滤无效数据,增量同步按时间字段筛选新增数据
    # df = df.filter(df["create_time"] > "上次同步的最大时间戳")

    # 步骤2:写入Elasticsearch
    df.write.format("es") \
        # 写入模式:overwrite覆盖全量/append追加增量/ignore已存在则跳过/error已存在则报错
        .mode("append") \
        # 可选:指定字段作为ES文档的_id,比如上面重命名后的mongo_id
        # .option("es.mapping.id", "mongo_id") \
        # ES7+直接填索引名即可,低版本需要填索引/类型
        .save("目标ES索引名")

    spark.stop()

核心参数说明

  • 连接器版本需要和你的Spark版本、Scala版本、ES版本匹配,否则会出现类找不到报错
  • spark.es.nodes.wan.only:如果访问的是托管云ES服务、或者ES和Spark不在同一个内网,必须设为true,否则会因为节点自动发现失败报错
  • 全量同步场景写入模式用overwrite,增量同步场景用append即可

任务提交命令示例

spark-submit --packages org.mongodb.spark:mongo-spark-connector_2.12:3.0.1,org.elasticsearch:elasticsearch-spark-30_2.12:7.17.0 你的脚本文件名.py

上述依赖版本适配Spark3.0+、Scala2.12、ES7.x版本,可根据自己的环境版本替换对应包版本号

内容的提问来源于stack exchange,提问作者sam.marhaendra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 15:39:03