【问题标题】:unable to send data to MongoDB using spark strucutred streaming无法使用 Spark 结构化流向 MongoDB 发送数据
【发布时间】:2020-12-06 11:31:02
【问题描述】:

我按照这个Unable to send data to MongoDB using Kafka-Spark Structured Streaming 将数据从 spark 结构化流发送到 mongoDB,我成功地实现了它,但是有一个问题。 比如当函数

override def process(record: Row): Unit = {

    val doc: Document = Document(record.prettyJson.trim)
    // lazy opening of MongoDB connection


    ensureMongoDBConnection()
    val result = collection.insertOne(doc)
    if (messageCountAccum != null)
      messageCountAccum.add(1)
  }

代码执行没有任何问题,但没有数据发送到 MongoDB

但是如果我添加这样的打印语句

override def process(record: Row): Unit = {
    val doc: Document = Document(record.prettyJson.trim)

    // lazy opening of MongoDB connection


    ensureMongoDBConnection()
    val result = collection.insertOne(doc)
    result.foreach(println) //print statement
    if (messageCountAccum != null)
      messageCountAccum.add(1)
  }

数据正在被插入 MongoDB

不知道为什么???

【问题讨论】:

    标签: mongodb scala spark-streaming spark-structured-streaming


    【解决方案1】:

    foreach 初始化写入器接收器。如果没有 foreach,则永远不会计算您的数据框。

    Try this :
    
    val df = // your df here
    df.map(r => process(r))
    df.count()
    

    【讨论】:

      猜你喜欢
      • 2019-09-11
      • 2018-08-20
      • 1970-01-01
      • 2019-08-16
      • 2019-08-24
      • 1970-01-01
      • 2017-07-12
      • 1970-01-01
      • 2017-05-04
      相关资源
      最近更新 更多