Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

382
Visualizações
El uso de Flux.buffer de reactor para el trabajo por lotes solo funciona para un solo artículo

Estoy tratando de usar Flux.buffer() para agrupar cargas desde una base de datos.

El caso de uso es que la carga de registros desde una base de datos puede ser "a ráfagas", y me gustaría presentar un pequeño búfer para agrupar las cargas cuando sea posible.

Mi enfoque conceptual ha sido usar alguna forma de procesador, publicar en su receptor, dejar que se almacene en búfer y luego suscribirse y filtrar para obtener el resultado que quiero.

Probé múltiples enfoques diferentes (diferentes tipos de procesadores, creando el Mono filtrado de diferentes maneras).

A continuación se muestra a dónde he llegado hasta ahora, en gran parte por tropezar.

Actualmente, esto devuelve un solo resultado, pero las llamadas posteriores se eliminan (aunque no estoy seguro de dónde).

 class BatchLoadingRepository { // I've tried all manner of different processors here. I'm unsure if // TopicProcessor is the correct one to use. private val bufferPublisher = TopicProcessor.create<String>() private val resultsStream = bufferPublisher .bufferTimeout(50, Duration.ofMillis(50)) // I'm unsure if concatMapIterable is the correct operator here, // but it seems to work. // I'm really trying to turn the List<MyEntity> // into a stream of MyEntity, published on the Flux<> .concatMapIterable { requestedIds -> // this is a Spring Data repository. It returns List<MyEntity> repository.findAllById(requestedIds) } // Multiple callers will invoke this method, and then subscribe to receive // their entity back. fun findByIdAsync(id: String): Mono<MyEntity> { // Is there a potential race condition here, caused by a result // on the resultsStream, before I've subscribed? return Mono.create<MyEntity> { sink -> bufferPublisher.sink().next(id) resultsStream.filter { it.id == id } .subscribe { next -> sink.success(next) } } } }
over 4 years ago · Santiago Trujillo
1 Respostas
Responde à pergunta

0

Hola, estaba probando tu código y creo que la mejor manera es usar EmitterProcessor shared. Hice una prueba con emitterProcessor y parece funcionar.

 Flux<String> fluxi; EmitterProcessor emitterProcessor; @Override public void run(String... args) throws Exception { emitterProcessor = EmitterProcessor.create(); fluxi = emitterProcessor.share().bufferTimeout(500, Duration.ofMillis(500)) .concatMapIterable(o -> o); Flux.range(0,1000) .flatMap(integer -> findByIdAsync(integer.toString())) .map(s -> { System.out.println(s); return s; }).subscribe(); } private Mono<String> findByIdAsync(String id) { return Mono.create(monoSink -> { fluxi.filter(s -> s == id).subscribe(value -> monoSink.success(value)); emitterProcessor.onNext(id); }); }
over 4 years ago · Santiago Trujillo Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda