- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我想通过 rsocket 连接两个应用程序。第一个是用 GO 编写的,第二个是用 Kotlin 编写的。我想实现客户端发送数据流和服务器发送确认响应的连接。
问题在于等待所有元素,如果服务器不 BlockOnLast(ctx),则读取整个流,但在所有条目到达之前发送响应。如果添加了BlockOnLast(ctx),则服务器(GoLang)会卡住。
我还用 Kotlin 编写了客户端,在这种情况下,整个通信工作得很好。
有人可以帮忙吗?
GO 服务器:
package main
import (
"context"
"github.com/golang/protobuf/proto"
"github.com/rsocket/rsocket-go"
"github.com/rsocket/rsocket-go/payload"
"github.com/rsocket/rsocket-go/rx"
"github.com/rsocket/rsocket-go/rx/flux"
"rsocket-go-rpc-test/proto"
)
func main() {
addr := "tcp://127.0.0.1:8081"
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
err := rsocket.Receive().
Fragment(1024).
Resume().
Acceptor(func(setup payload.SetupPayload, sendingSocket rsocket.CloseableRSocket) (rsocket.RSocket, error) {
return rsocket.NewAbstractSocket(
rsocket.RequestChannel(func(payloads rx.Publisher) flux.Flux {
println("START")
payloads.(flux.Flux).
DoOnNext(func(input payload.Payload) {
chunk := &pl_dwojciechowski_proto.Chunk{}
proto.Unmarshal(input.Data(), chunk)
println(string(chunk.Content))
}).BlockLast(ctx)
return flux.Create(func(i context.Context, sink flux.Sink) {
status, _ := proto.Marshal(&pl_dwojciechowski_proto.UploadStatus{
Message: "OK",
Code: 0,
})
sink.Next(payload.New(status, make([]byte, 1)))
sink.Complete()
println("SENT")
})
}),
), nil
}).
Transport(addr).
Serve(ctx)
panic(err)
}
Kotlin 客户端:
private fun clientCall() {
val rSocket = RSocketFactory.connect().transport(TcpClientTransport.create(8081)).start().block()
val client = FileServiceClient(rSocket)
val requests: Flux<Chunk> = Flux.range(1, 10)
.map { i: Int -> "sending -> $i" }
.map<Chunk> {
Chunk.newBuilder()
.setContent(ByteString.copyFrom(it.toByteArray())).build()
}
val response = client.send(requests).block() ?: throw Exception("")
rSocket.dispose()
System.out.println(response.message)
}
以及用 Kotlin 编写的 GO 的等效项:
val serviceServer = FileServiceServer(DefaultService(), Optional.empty(), Optional.empty())
val closeableChannel = RSocketFactory.receive()
.acceptor { setup: ConnectionSetupPayload?, sendingSocket: RSocket? ->
Mono.just(
RequestHandlingRSocket(serviceServer)
)
}
.transport(TcpServerTransport.create(8081))
.start()
.block()
closeableChannel.onClose().block()
class DefaultService : FileService {
override fun send(messages: Publisher<Service.Chunk>?, metadata: ByteBuf?): Mono<Service.UploadStatus> {
return Flux.from(messages)
.windowTimeout(10, Duration.ofSeconds(500))
.flatMap(Function.identity())
.doOnNext { println(it.content.toStringUtf8()) }
.then(Mono.just(Service.UploadStatus.newBuilder().setCode(Service.UploadStatusCode.Ok).setMessage("test").build()))
}
}
服务器输出:
START
sending -> 1
最佳答案
解决方案如下:
package main
import (
"context"
"github.com/golang/protobuf/proto"
"github.com/rsocket/rsocket-go"
"github.com/rsocket/rsocket-go/payload"
"github.com/rsocket/rsocket-go/rx"
"github.com/rsocket/rsocket-go/rx/flux"
"rsocket-go-rpc-test/proto"
)
type TestService struct {
totals int
pl_dwojciechowski_proto.FileService
}
var statusOK = &pl_dwojciechowski_proto.UploadStatus{
Message: "code",
Code: pl_dwojciechowski_proto.UploadStatusCode_Ok,
}
var statusErr = &pl_dwojciechowski_proto.UploadStatus{
Message: "code",
Code: pl_dwojciechowski_proto.UploadStatusCode_Failed,
}
func main() {
addr := "tcp://127.0.0.1:8081"
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
err := rsocket.Receive().
Fragment(1024).
Acceptor(func(setup payload.SetupPayload, sendingSocket rsocket.CloseableRSocket) (rsocket.RSocket, error) {
return rsocket.NewAbstractSocket(
rsocket.RequestChannel(func(msgs rx.Publisher) flux.Flux {
dataReceivedChan := make(chan bool, 1)
toChan, _ := flux.Clone(msgs).
DoOnError(func(e error) {
dataReceivedChan <- false
}).
DoOnComplete(func() {
dataReceivedChan <- true
}).
ToChan(ctx, 1)
fluxResponse := flux.Create(func(ctx context.Context, s flux.Sink) {
gluedContent := make([]byte, 1024)
for c := range toChan {
chunk := pl_dwojciechowski_proto.Chunk{}
_ = chunk.XXX_Unmarshal(c.Data())
gluedContent = append(gluedContent, chunk.Content...)
}
if <-dataReceivedChan {
marshal, _ := proto.Marshal(statusOK)
s.Next(payload.New(marshal, nil))
s.Complete()
} else {
marshal, _ := proto.Marshal(statusErr)
s.Next(payload.New(marshal, nil))
s.Complete()
}
})
return fluxResponse
}),
), nil
}).
Transport(addr).
Serve(ctx)
panic(err)
}
关于java - 等待最后一个元素时阻塞 Flux,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/59735193/
我试图让脚本暂停大约 1 秒,然后继续执行脚本,但我似乎无法弄清楚如何做。这是我的代码: function hello() { alert("Hi!") //I need about a 1
wait() 和 wait(timeout) 之间有什么区别。无论如何 wait() 需要等待通知调用,但为什么我们有 wait(timeout)? 那么 sleep(timeout) 和 wait(
我需要做什么: 我有一个带有文件输入和隐藏文本输入的上传表单。用户上传图像,图像被操作,然后发送到远程服务器进行处理,这需要几秒钟,然后远程服务器将最终的图像发送回家庭服务器,并保存在新文件夹中。 J
大家好,我正在使用 Visual C++ 2010,尝试使用 Winsock 编写服务器/客户端应用程序...我不确定为什么,但有时服务器会在 listen() 函数处等待,有时会在 accept 处
任务描述 我为我的 Angular 应用程序实现了 CRSF 保护。服务器检查 crsf token 是否位于请求的 header “X-CSRF-TOKEN”中。如果不是,它会发送一个 HTTP 响
我想做这个例子https://stackoverflow.com/a/33585993/1973680同步。 这是正确的实现方式吗? let times= async (n,f)=>{
我如何将 while 循环延迟到 1 秒间隔,而不会将其运行的整个代码/计算机的速度减慢到一秒延迟(只是一个小循环)。 最佳答案 Thread.sleep(1000); // do nothing f
我知道这是一个重复的问题。但是我无法通过解释来理解。我想用一个很好的例子来清楚地理解它。任何人都可以帮忙吗。 “为什么我们从同步上下文中调用 wait()、notify() 方法”。 最佳答案 当我们
我有一个 click 事件,该事件是第一次从另一个地方自动触发的。我的问题是它运行得太快,因为所需的变量仍在由 Flash 和 Web 服务定义。所以现在我有: (function ($) {
我有如下功能 function async populateInventories(custID){ this.inventories = await this.inventoryServic
我一直对“然后”不被等待的行为感到困扰,我明白其原因。然而,我仍然需要绕过它。这是我的用例。 doWork(family) { return doWork1(family)
我想我理解异步背后的想法,返回一个Future,但是我不清楚异步在一个非常基本的层面上如何表现。据我了解,它不会自动在程序中创建异步行为。例如: import 'dart:async'; main()
我正在制作一个使用异步的Flutter应用程序,但它的工作方式不像我对它的了解。所以我对异步和在 Dart 中等待有一些疑问。这是一个例子: Future someFunction() async {
我在 main.tf 中创建资源组和 vNet,并在同一文件中引用模块。问题是,模块无法从模块访问这些资源。相关代码(删除了大部分代码,只留下相关部分): main.tf: module "worke
我的代码的问题是,当代码第一次运行时,我试图获取的 dom 元素并不总是存在,如果它不存在,那么永远不会做出 promise 。 我是否可以等到 promise 做出后再尝试实现它? 我希望我的最后一
所以,过去几天我一直在研究这段代码,并尝试实现回调/等待/任何需要的东西,但没有成功。 问题是,我如何等待响应,直到我得到两个函数的回调? (以及我将如何实现) 简而言之,我想做的是: POST 发生
谁能帮我理解这一点吗? 如果我们有一个类: public class Sample{ public synchronized method1(){ //Line1 .... wait();
这是我编写的代码,用于测试 wait() 和 notify() 的工作。现在我有很多疑问。 class A extends Thread { public void run() { try
我有以下代码由于语法错误而无法运行(在异步函数外等待) 如何使用 await 定义变量并将其导出? 当我这样定义一个变量并从其他文件导入它时,该变量是只创建一次(第一次读取文件时?)还是每次导入时都创
一个简单的线程程序,其中写入器将内容放入堆栈,读取器从堆栈中弹出。 java.util.Stack; import java.util.concurrent.ExecutorService; impo
我是一名优秀的程序员,十分优秀!