【问题标题】:how retrieve and add new document to MongoDb collection throw com.mongodb.reactivestreams.client.MongoClient如何检索新文档并将其添加到 MongoDb 集合中抛出 com.mongodb.reactivestreams.client.MongoClient
【发布时间】:2021-02-26 15:17:41
【问题描述】:

上下文:我编写了一个接收简单消息的 Kafka Consumer,我想使用 com.mongodb.reactivestreams.client.MongoClient 将其插入 MongoDb。虽然我理解我的问题是关于如何正确使用 MongoClient 让我通知我的堆栈:我的堆栈是 Micronaut + MongoDb reactive + Kotlin。

免责声明:如果有人在 java 中提供答案,我可以将其翻译成 Kotlin。您可以忽略下面的 Kafka 部分,因为它按预期工作。

这是我的代码

package com.mybank.consumer

import com.mongodb.reactivestreams.client.MongoClient
import com.mongodb.reactivestreams.client.MongoCollection
import com.mongodb.reactivestreams.client.MongoDatabase
import io.micronaut.configuration.kafka.annotation.KafkaKey
import io.micronaut.configuration.kafka.annotation.KafkaListener
import io.micronaut.configuration.kafka.annotation.OffsetReset
import io.micronaut.configuration.kafka.annotation.Topic
import org.bson.Document
import org.reactivestreams.Publisher
import javax.inject.Inject


@KafkaListener(offsetReset = OffsetReset.EARLIEST)
class DebitConsumer {

    @Inject
    //@Named("another")
    var mongoClient: MongoClient? = null


    @Topic("debit")
    fun receive(@KafkaKey key: String, name: String) {


        println("Account - $name by $key")

        
        var mongoDb : MongoDatabase? = mongoClient?.getDatabase("account")
        var mongoCollection: MongoCollection<Document>? = mongoDb?.getCollection("account_collection")
        var mongoDocument: Publisher<Document>? = mongoCollection?.find()?.first()
        print(mongoDocument.toString())

        //println(mongoClient?.getDatabase("account")?.getCollection("account_collection")?.find()?.first())
        //val mongoClientClient: MongoDatabase  = mongoClient.getDatabase("account")
        //println(mongoClient.getDatabase("account").getCollection("account_collection").find({ "size.h": { $lt: 15 } })
        //println(mongoClient.getDatabase("account").getCollection("account_collection").find("1").toString())


    }
}

嗯,上面的代码是我得到的最接近的。它没有提示任何错误。正在打印

com.mongodb.reactivestreams.client.internal.Publishers$$Lambda$618/0x0000000800525840@437ec11

我想这证明代码正确连接到数据库,但我希望打印第一个文档。

共有三个文件:

我的最终目标是将我从 Kafka Listener 收到的消息插入到 MongoDb。任何线索将不胜感激。

整个代码可以在git hub找到

***在苏珊的问题后编辑

这是打印的内容

var mongoDocument = mongoCollection?.find()?.first()
print(mongoDocument.toString())

【问题讨论】:

  • “但我希望打印第一个文档”- 我认为您的代码是为打印发布者而不是文档而编写的。
  • 这打印什么? var mongoDocument = mongoCollection?.find()?.first()
  • @SusanMustafa 我在打印上方添加了

标签: mongodb kotlin micronaut reactive-mongo-java


【解决方案1】:

看起来您正在为 mongodb 使用响应式流。您使用响应式流是有原因的吗?

您得到的结果是“Publisher”类型。您需要使用方法 subscribe() 来获取文档。

请参阅 Publisher 上的文档。

http://www.howsoftworks.net/reacstre/1.0.2/Publisher

如果您不想使用响应式:关于如何/如何在 Kotlin 中使用 mongodb 的绝佳示例。

https://kb.objectrocket.com/mongo-db/retrieve-mongodb-document-using-kotlin-1180

--- 使用 MongoDB、Reactive Streams、Publisher 的类似 StackOverlow。

how save document to MongoDb with com.mongodb.reactivestreams.client

===============已编辑==============

Publisher<Document> publisher = collection.find().first();

subscriber = new PrintDocumentSubscriber();
publisher.subscribe(subscriber); //publisher.subscribe(subscriber)
subscriber.await();

The example will print the following document:

{ "_id" : { "$oid" : "551582c558c7b4fbacf16735" },
  "name" : "MongoDB", "type" : "database", "count" : 1,
}

如果你想要非阻塞,这样做:

publisher.subscribe(new PrintDocumentSubscriber());  //without await

http://mongodb.github.io/mongo-java-driver-reactivestreams/1.4/javadoc/tour/SubscriberHelpers.PrintDocumentSubscriber.html

http://mongodb.github.io/mongo-java-driver-reactivestreams/1.6/getting-started/quick-tour/

【讨论】:

  • 是的,这是有原因的。我想建立一个依赖于 Micronaut 的端到端没有阻塞的反应式堆栈。打个比方,假设我使用的是 Spring WebFlux/Netty,我会使用 spring-boot-starter-data-mongodb-reactive。实际上,我今年使用 ElasticSearch 编写了两个基于 Spring WebFlux/Spring Data Reactive 的解决方案。我距离成为 Reactive pradigma 的专家还很遥远,但据我所知,如果驱动程序到数据库被阻塞,我们就会失去响应式设计的好处。因为我想使用 Micronaut 并避免使用 Spring,所以我正在尝试 com.mongodb.reactivestreams.client.MongoClient
  • @JimC 你比我更了解reactive.streams :),我希望你能找到一些帮助的链接。你肯定需要订阅你的发布者。
  • 如果您可以编辑您的答案并添加一个如何订阅我的出版商的建议,我将不胜感激。使用 Spring Data 似乎更容易,可能是因为它隐藏了一些特定的 Reactive 特性并使它看起来像一个常见的 CRUD(讨论为什么我不简单地跳回 Spring 超出了这个问题)
  • 完成 @JimC mongodb github 文档非常有帮助。慢慢阅读,希望一切顺利
  • @JimC 对不起,我实际上没有。我希望如果一切顺利,你会给出一个答案来帮助未来的开发者
猜你喜欢
  • 2021-04-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-01-05
  • 1970-01-01
  • 2021-10-09
相关资源
最近更新 更多