【问题标题】:Send continuous query data from apache ignite to rsocket with request-stream使用请求流将连续查询数据从 apache ignite 发送到 rsocket
【发布时间】:2021-04-14 17:34:01
【问题描述】:

我正在尝试在 ignite 连续查询中接收数据,并通过 rsocket 请求流向客户端发送流数据。连续查询的代码块是这样的:

 ContinuousQuery<String, BinaryObject> query = new ContinuousQuery<>();

        query.setLocalListener(new CacheEntryUpdatedListener<String, BinaryObject>() {
            @Override
            public void onUpdated(Iterable<CacheEntryEvent<? extends String, ? extends BinaryObject>> events)
                    throws CacheEntryListenerException {
                // react to the update events here
                events.forEach(cacheEntryEvent -> {
                    SampleResponse sampleResponse= createResponse(cacheEntryEvent.getValue());
              });
            }
        });
        igniteCache.query(query);

此代码在 Spring Boot 服务中。我想知道将 sampleResponse 发送到控制器类中的 rsocket,如下所示:

@MessageMapping("marketData.stream")
    public Flux<SampleResponse> responseStream(@Header("jwt") String jwt, SampleRequest request) {
// the method which contains the continuous query    
 return Flux.push(service.processStream(request));
}

如何通知 rsocket 发送在 ignite 连续查询的 onUpdate 方法中创建的消息?

【问题讨论】:

    标签: java ignite rsocket


    【解决方案1】:

    一般模式是Flux.create

        val f: Flux<String> = Flux.create { 
          it.onCancel { 
            // cancel subscription
          }
          
          it.next("a")
          it.next("b")
          it.next("c")
          it.complete()
        }
    

    对你来说很有可能

        val f: Flux<String> = Flux.create { 
          ContinuousQuery<String, BinaryObject> query = new ContinuousQuery<>();
    
    
            query.setLocalListener(new CacheEntryUpdatedListener<String, BinaryObject>() {
                @Override
                public void onUpdated(Iterable<CacheEntryEvent<? extends String, ? extends BinaryObject>> events)
                        throws CacheEntryListenerException {
                    // react to the update events here
                    events.forEach(cacheEntryEvent -> {
                        SampleResponse sampleResponse= createResponse(cacheEntryEvent.getValue());
                  });
                }
            });
            igniteCache.query(query);
    
            
          it.onCancel { 
            // cancel somehow
            xxx.cancel();
          }
        }
    

    但理想情况下,您应该有一些 ignite-reactor 或 ignite-jdk9flowable 绑定并使用它来代替。还是基于 Kotlin Flow 的东西?

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2023-01-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-10-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多