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

使用Spark-MongoDB Connector Scala时出现MongoTimeoutException的原因排查

嘿,我一眼就看到问题所在了!你遇到的MongoTimeoutException根本原因是连接错了端口,还有几个认证相关的小细节没处理好,我给你一步步拆解:

核心问题:端口搞混了

你的错误日志里显示连接的是localhost:8081,但这个端口是**mongo-express(MongoDB的Web管理界面)**的端口,不是MongoDB服务本身的端口!

看你的docker-compose配置:

  • mongo服务默认使用MongoDB的标准端口27017,只是你没把这个端口映射到宿主机
  • mongo-express才是用8081端口提供Web管理界面的服务

另外,你部署MongoDB的时候设置了用户名和密码,但代码里的连接URI完全没带认证信息,这也会导致连接失败。

修复步骤

1. 给MongoDB服务映射端口

修改你的docker-compose.yml,给mongo服务加上端口映射,让宿主机能直接访问容器里的MongoDB:

version: '3.3'
services:
  kafka:
    image: spotify/kafka
    ports:
      - "9092:9092"
    environment:
      - ADVERTISED_HOST=localhost
  mongo:
    image: mongo
    restart: always
    ports:
      - "27017:27017"  # 新增这行,把容器内的27017端口映射到宿主机
    environment:
      MONGO_INITDB_ROOT_USERNAME: user
      MONGO_INITDB_ROOT_PASSWORD: password
  mongo-express:
    image: mongo-express
    restart: always
    ports:
      - 8081:8081
    environment:
      ME_CONFIG_MONGODB_ADMINUSERNAME: user
      ME_CONFIG_MONGODB_ADMINPASSWORD: password

然后重启docker-compose服务:

docker-compose down && docker-compose up -d

2. 修正Scala代码里的连接URI

把所有的mongodb://localhost:8081/admin.new_col改成带用户名密码、正确端口的URI,还要指定认证源(因为你用的是admin数据库的root用户):

object SparkTest {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .master("local")
      .appName("MongoSparkConnectorIntro")
      // 修正后的URI:端口27017,带上用户名密码,指定authSource=admin
      .config("spark.mongodb.input.uri", "mongodb://user:password@localhost:27017/admin.new_col?authSource=admin")
      .config("spark.mongodb.output.uri", "mongodb://user:password@localhost:27017/admin.new_col?authSource=admin")
      .getOrCreate()

    val sparkContext = spark.sparkContext
    sparkContext.setLogLevel("ERROR")

    val writeConfig = WriteConfig(Map("uri" -> "mongodb://user:password@localhost:27017/admin.new_col?authSource=admin"))
    val documents = sparkContext.parallelize((1 to 10).map(i => Document.parse(s"{test: $i}")))
    MongoSpark.save(documents, writeConfig)
  }
}

3. 验证连接是否正常(可选但推荐)

运行代码前,你可以用MongoDB客户端测试连接:

mongosh "mongodb://user:password@localhost:27017/admin?authSource=admin"

如果能成功进入MongoDB的shell界面,说明连接配置没问题,再运行你的Scala代码就不会报超时错误了。

总结

简单来说,你就是把MongoDB服务的端口和它的Web管理工具端口搞混了,再加上没带认证信息,才导致Spark连不上MongoDB。按上面的步骤调整后就能解决问题啦!

内容的提问来源于stack exchange,提问作者Cassie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:50:46