【问题标题】:Quarkus SSE Redis subscribeQuarkus SSE Redis 订阅
【发布时间】:2021-03-26 19:52:50
【问题描述】:

我喜欢用 quarkus 中 redis.subscribe 的响应做一个 SSE。

我有一个来自 quarkus-quickstart 的简单 SSE 示例

 @GET
  @Produces(MediaType.SERVER_SENT_EVENTS)
  @SseElementType(MediaType.TEXT_PLAIN)
  @Path("{name}/streaming")
  public Multi<String> greeting(@org.jboss.resteasy.annotations.jaxrs.PathParam String name) {
    return Multi.createFrom().publisher(vertx.periodicStream(2000).toMulti())
        .map(l -> String.format("Hello %s! (%s)%n", name, new Date()));
  }

效果很好,每 2 秒我都会在我的网络浏览器中收到 Hello ....

现在我尝试订阅 Redis,所以我应该会收到来自 Redis 的消息。

Redis 示例:

(cmd window 1)
SUBSCRIBE message-channel
Reading messages... (press Ctrl-C to quit)
1) "subscribe"
2) "message-channel"
3) (integer) 1

(cmd window 2)
PUBLISH  message-channel HelloWorld
(integer) 1

(cmd window 1)
1) "message"
2) "message-channel"
3) "HelloWorld"

现在我用 quarkus SSE 试试这个:

  @Inject
  ReactiveRedisClient reactiveRedisClient;

 @GET
  @Produces(MediaType.SERVER_SENT_EVENTS)
  @SseElementType(MediaType.TEXT_PLAIN)
  @Path("sse/redissse")
  public Multi<String> redissse() {
    List<String> subscriberList = new ArrayList();
    subscriberList.add("message-channel");

    return reactiveRedisClient.subscribe(subscriberList)
        .onItem().transformToMulti(keys -> Multi.createFrom().iterable(keys))
        .onItem().castTo(String.class);
  }

我收到的是一个例外:

WARNING [io.ver.red.cli.imp.RedisConnectionImpl] (vert.x-eventloop-thread-0) No handler waiting for message: [subscribe, message-channel, 1]

有人可以支持我吗? 有一个简单的例子吗? 我对此一无所知,我无法通过“订阅”发布接收 Redis 消息。

任何建议...

【问题讨论】:

    标签: java redis publish-subscribe server-sent-events quarkus


    【解决方案1】:

    我没有使用过 Redis pub-sub,但我确实使用过 Redis 流,我必须做的是这样的:

    `

    return Multi.createBy().repeating()
        .supplier(() -> this.reactiveRedisClient.subscribe(subscriberList)
                            .onItem().transformToMulti(keys -> Multi.createFrom().iterable(keys))
                            .onItem().castTo(String.class))
            .indefinitely()
            .onItem().disjoint();
    

    `

    我猜由于 pub-sub 是非阻塞的,它运行一次然后它不会等到另一条消息到达。您必须以反应方式实现自己的 while(true) 循环。

    【讨论】:

      【解决方案2】:

      现在我执行以下操作:

        @Inject
        @RedisClientName("second")
        RedisClient redisClient2;
      
      void onStart(@Observes StartupEvent ev) throws IOException {
        this.redisClient2.subscribe(List.of("message-channel"));
      }
      
      
        @GET
        @Produces(MediaType.SERVER_SENT_EVENTS)
        @SseElementType(MediaType.TEXT_PLAIN)
        @Path("/redis/subscribe")
        public Publisher<String> subscribechannel(){
           return eventBus.<String>consumer("io.vertx.redis.message-channel").toPublisherBuilder()
              .map(Message::body)
              .buildRs();
        }
      

      现在它可以工作了,但是如果我从多个浏览器进行 SSE,他们会共享事件。因此,每个消费者(浏览器)之后只有一个收到事件。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-05-22
        • 2015-09-04
        • 1970-01-01
        • 2013-10-11
        • 2021-04-12
        相关资源
        最近更新 更多