【问题标题】:How to write an RSocket client in JavaScript如何用 JavaScript 编写 RSocket 客户端
【发布时间】:2019-09-18 12:54:22
【问题描述】:

我尝试用 Java 实现一个 RSocket 服务器,用 JavaScript 实现一个客户端,但我无法调用后端中的任何方法。

Java 服务器

public final class RawServer {
    public static void main(String[] args) {
        RSocketFactory.receive()
                .acceptor((setup, sendingSocket) -> Mono.just(new DefaultSimpleService()))
                .transport(WebsocketServerTransport.create("localhost", 8801))
                .start()
                .block()
                .onClose()
                .block();
    }

    private static final class DefaultSimpleService extends AbstractRSocket {
        private ObjectMapper jsonMapper = new ObjectMapper();

        @Override
        public Flux<Payload> requestStream(Payload payload) {
            return Mono.just(payload.getDataUtf8())
                    .map(json -> {
                        try {
                            return jsonMapper.readValue(json, Message.class);
                        } catch (IOException e) {
                            e.printStackTrace();
                            return null;
                        }
                    })
                    .doOnNext(msg -> System.out.println("got message " + msg.message))
                    .flatMapMany(msg -> Flux.range(0, 5)
                            .map(count -> msg.message + " #" + count))
                    .map(message -> DefaultPayload.create(message));
        }
    }
}

public class Message {

    public final String message;

    @JsonCreator
    public Message(@JsonProperty("message") String message) {
        this.message = message;
    }
}

JavaScript 客户端

    import { RSocketClient, JsonSerializers } from "rsocket-core";
    import RSocketWebSocketClient from "rsocket-websocket-client";

    const transport = new RSocketWebSocketClient({
        url: "ws://localhost:8801"
      });

      const client = new RSocketClient({
        // send/receive JSON objects instead of strings/buffers
        serializers: JsonSerializers,
        setup: {
          // ms btw sending keepalive to server
          keepAlive: 60000,
          // ms timeout if no keepalive response
          lifetime: 180000,
          // format of `data`
          dataMimeType: "application/json",
          // format of `metadata`
          metadataMimeType: "application/json"
        },
        transport
      });
      client.connect().subscribe({
        onComplete: socket => {
          socket.requestStream({
            data: { message: "hello from javascript!" },
            metadata: null
          });
        },
        onError: error => {
          console.log("got error");
          console.error(error);
        },
        onSubscribe: cancel => {
          /* call cancel() to abort */
          console.log("subscribe!");
          console.log(cancel);
          // cancel.cancel();
        }
      });

WebSocket 连接似乎已建立,但没有消息推送到服务器。我该怎么做?

我还用 Java 实现了客户端,它工作得很好。我找到的 JavaScript 示例是 https://github.com/rsocket/rsocket-js/blob/master/docs/01-client-configuration.md,但我无法使其工作。

【问题讨论】:

    标签: javascript websocket reactive-streams rsocket


    【解决方案1】:

    更多关于 RSocket 的例子你可以访问我的个人博客http://kojotdev.com/2019/09/rsocket-examples-java-javascript-spring-webflux/

    好的,我想通了。首先,我们需要修复我们的服务器以返回正确的 JSON 对象。

    @Override
    public Flux<Payload> requestStream(Payload payload) {
        log.info("got requestStream in Server");
        log.info(payload.getDataUtf8());
        return Mono.just(payload.getDataUtf8())
                .map(payloadString -> MessageMapper.jsonToMessage(payloadString))
                .flatMapMany(msg -> Flux.range(0, 5)
                        .map(count -> msg.message + " | requestStream from Server #" + count)
                        .map(responseText -> new Message(responseText))
                        .map(responseMessage -> MessageMapper.messageToJson(responseMessage)))
                .map(message -> DefaultPayload.create(message));
    }
    

    然后,在我们的 JavaScript 客户端中,我们需要将 socket.requestStream 方法更改为:

    socket
      .requestStream({
        data: { message: "request - stream from javascript!" },
        metadata: ""
      })
      .subscribe({
        onComplete: () => console.log("requestStream done"),
        onError: error => {
          console.log("got error with requestStream");
          console.error(error);
        },
        onNext: value => {
          // console.log("got next value in requestStream..");
          console.log(value.data);
        },
        // Nothing happens until `request(n)` is called
        onSubscribe: sub => {
          console.log("subscribe request Stream!");
          sub.request(7);
        }
      });
    

    其他一切都和前面的例子一样。 有用的链接:

    【讨论】:

    • 为什么如果我执行此代码两次它会告诉我客户端已经连接?
    • 因为您可能两次调用client.connect()。在client.connect()onComplete 事件中,将socket 存储到全局变量中。然后调用socket.requestResponse()等方法
    猜你喜欢
    • 2022-12-29
    • 1970-01-01
    • 2015-07-25
    • 2015-06-15
    • 2012-06-01
    • 2015-10-23
    • 1970-01-01
    • 2022-10-24
    • 1970-01-01
    相关资源
    最近更新 更多