【发布时间】: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