Java vertx webclient Stream с сервера с кусочками сообщенийJAVA

Программисты JAVA общаются здесь
Anonymous
Java vertx webclient Stream с сервера с кусочками сообщений

Сообщение Anonymous »

У меня есть сервер Vertx, который отправляет на события для всех стримеров информацию.
Использование консоли, которая работает так, как желает.

Код: Выделить всё

val webClientOpt = WebClientOptions()
.setKeepAlive(true)
.setUserAgent("Client/2.0")
.setFollowRedirects(true)
.setShared(true)
< /code>
Client Call < /p>
client
.get(port, host, UriTemplate.of(s"${path}event-bus/"))
.putHeader("content-type", "application/json")
.bearerTokenAuthentication(UserBuffer.loggedInUser().jwtToken())
.as(BodyCodec.pipe(writeBuffer)).send()
< /code>
Теперь writebuffer и читатель < /p>
val writeBuffer = ReactiveWriteStream.writeStream[Buffer](vertx)
val readStream = ReactiveReadStream.readStream[Buffer]()
< /code>
также материал, чтобы получить данные < /p>
readStream.handler(j => {
println("CONSUMING!!!!")
println(j.toString("UTF-8"))
})
writeBuffer.subscribe(readStream)
< /code>
Я знаю, что должен использовать BodyCodec.pipe, но я думаю, что вот моя проблема. Я думаю, что я не использую его правильно.new WriteStream[Buffer]() {
override def write(buffer: Buffer): io.vertx.core.Future[Void] = {
println(buffer.toString())
Future.successful(null).asVertx
}
override def end(): io.vertx.core.Future[Void] = {
Future.successful(null).asVertx
}
override def setWriteQueueMaxSize(maxSize: Int): WriteStream[Buffer] = this
override def writeQueueFull(): Boolean = false
override def drainHandler(handler: Handler[Void]): WriteStream[Buffer] = this
override def exceptionHandler(handler: Handler[Throwable]): WriteStream[Buffer] = this
}
Это то, что я искал.

Подробнее здесь: https://stackoverflow.com/questions/797 ... d-messages

Вернуться в «JAVA»