如何使用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
相关产品推荐
相关产品推荐

