AskOverflow.Dev

AskOverflow.Dev Logo AskOverflow.Dev Logo

AskOverflow.Dev Navigation

  • Início
  • system&network
  • Ubuntu
  • Unix
  • DBA
  • Computer
  • Coding
  • LangChain

Mobile menu

Close
  • Início
  • system&network
    • Recentes
    • Highest score
    • tags
  • Ubuntu
    • Recentes
    • Highest score
    • tags
  • Unix
    • Recentes
    • tags
  • DBA
    • Recentes
    • tags
  • Computer
    • Recentes
    • tags
  • Coding
    • Recentes
    • tags
Início / coding / Perguntas / 77674891
Accepted
MartinHH
MartinHH
Asked: 2023-12-17 23:11:40 +0800 CST2023-12-17 23:11:40 +0800 CST 2023-12-17 23:11:40 +0800 CST

fs2: Como fazer algo depois que o stream for iniciado ("doOnSubscribe")?

  • 772

Estou tentando usar uma API impura ("java") no contexto de um IO-App com efeito de gato. A API impura se parece com isto:

import  io.reactivex.Flowable
import java.util.concurrent.CompletableFuture

trait ImpureProducer[A] {
  /** the produced output - must be subscribed to before calling startProducing() */
  def output: Flowable[A]

  /** Makes the producer start publishing its output to any subscribers of `output`. */
  def startProducing(): CompletableFuture[Unit]
}

(É claro que existem mais métodos, incluindo stopProduzindo(), mas esses não são relevantes para minha pergunta.)

Minha abordagem (talvez ingênua) para adaptar essa API é a seguinte (aproveitando o fato de que Flowableé um org.reactivestreams.Publisher):

import cats.effect.IO
import fs2.Stream
import fs2.interop.reactivestreams.*

class AdaptedProducer[A](private val underlying: ImpureProducer[A]) {
  def output: Stream[IO, A] =
    underlying.output.toStreamBuffered(1)
  def startProducing: IO[Unit] = 
    IO.fromCompletableFuture(IO(underlying.startProducing()))
}

Minha pergunta é: como posso garantir que o output-stream esteja inscrito antes de avaliar startProducing?

Por exemplo, como eu poderia corrigir a seguinte tentativa de obter um IO do primeiro item produzido:

import cats.Parallel
import cats.effect.IO

def firstOutput[A](producer: AdaptedProducer[A]): IO[A] = {
  val firstOut: IO[A] = producer.output.take(1).compile.onlyOrError
  // this introduces a race condition: it is not ensured that the output-stream
  // will already be subscribed to when startProducing is evaluated.
  Parallel[IO].parProductL(firstOut)(producer.startProducing)
}
scala
  • 1 1 respostas
  • 25 Views

1 respostas

  • Voted
  1. Best Answer
    Luis Miguel Mejía Suárez
    2023-12-18T01:09:32+08:002023-12-18T01:09:32+08:00

    Isso pode ser feito usando a nova interoperabilidade Java flow .
    Observe que o inteop de fluxos reativos está obsoleto desde a adição dele flow. A API foi redesenhada para lidar com mais casos (como este) e a implementação foi refeita do zero para ser mais eficiente.
    Se a API que você está usando não tiver uma flowversão baseada, você poderá usar FlowAdapterspara envolvê-la.

    O código ficaria assim:

    import org.reactivestreams.FlowAdapters
    
    final class AdaptedProducer[A](underlying: ImpureProducer[A], chunkSize: Int) {
      val run: Stream[IO, A] =
        fs2.interop.flow.fromPublisher(chunkSize) { subscriber =>
           IO(
             underlying.output.subscribe(
                FlowAdapters.toFlowSubscriber(
                 subscriber
                )
             )
           ) >>
           IO.fromCompletableFuture(IO(underlying.startProducing()))
        }    
    }
    

    Isso garante que that startProducingseja chamado apenas quando você realmente começar a consumir Streame depois de ter chamado subscribe.

    Espero que isso ajude :D

    • 2

relate perguntas

  • Por que esta composição da camada ZIO não é compilada?

  • Como posso converter um Expr em uma Árvore no Scala 3?

  • Como combinar um TypeDef em uma anotação de macro Scala 3 corretamente?

  • Spark Scala mesclando várias colunas em uma única coluna

  • Argumento de método multitipo para varargs

Sidebar

Stats

  • Perguntas 205573
  • respostas 270741
  • best respostas 135370
  • utilizador 68524
  • Highest score
  • respostas
  • Marko Smith

    destaque o código em HTML usando <font color="#xxx">

    • 2 respostas
  • Marko Smith

    Por que a resolução de sobrecarga prefere std::nullptr_t a uma classe ao passar {}?

    • 1 respostas
  • Marko Smith

    Você pode usar uma lista de inicialização com chaves como argumento de modelo (padrão)?

    • 2 respostas
  • Marko Smith

    Por que as compreensões de lista criam uma função internamente?

    • 1 respostas
  • Marko Smith

    Estou tentando fazer o jogo pacman usando apenas o módulo Turtle Random e Math

    • 1 respostas
  • Marko Smith

    java.lang.NoSuchMethodError: 'void org.openqa.selenium.remote.http.ClientConfig.<init>(java.net.URI, java.time.Duration, java.time.Duratio

    • 3 respostas
  • Marko Smith

    Por que 'char -> int' é promoção, mas 'char -> short' é conversão (mas não promoção)?

    • 4 respostas
  • Marko Smith

    Por que o construtor de uma variável global não é chamado em uma biblioteca?

    • 1 respostas
  • Marko Smith

    Comportamento inconsistente de std::common_reference_with em tuplas. Qual é correto?

    • 1 respostas
  • Marko Smith

    Somente operações bit a bit para std::byte em C++ 17?

    • 1 respostas
  • Martin Hope
    fbrereto Por que a resolução de sobrecarga prefere std::nullptr_t a uma classe ao passar {}? 2023-12-21 00:31:04 +0800 CST
  • Martin Hope
    比尔盖子 Você pode usar uma lista de inicialização com chaves como argumento de modelo (padrão)? 2023-12-17 10:02:06 +0800 CST
  • Martin Hope
    Amir reza Riahi Por que as compreensões de lista criam uma função internamente? 2023-11-16 20:53:19 +0800 CST
  • Martin Hope
    Michael A formato fmt %H:%M:%S sem decimais 2023-11-11 01:13:05 +0800 CST
  • Martin Hope
    God I Hate Python std::views::filter do C++20 não filtrando a visualização corretamente 2023-08-27 18:40:35 +0800 CST
  • Martin Hope
    LiDa Cute Por que 'char -> int' é promoção, mas 'char -> short' é conversão (mas não promoção)? 2023-08-24 20:46:59 +0800 CST
  • Martin Hope
    jabaa Por que o construtor de uma variável global não é chamado em uma biblioteca? 2023-08-18 07:15:20 +0800 CST
  • Martin Hope
    Panagiotis Syskakis Comportamento inconsistente de std::common_reference_with em tuplas. Qual é correto? 2023-08-17 21:24:06 +0800 CST
  • Martin Hope
    Alex Guteniev Por que os compiladores perdem a vetorização aqui? 2023-08-17 18:58:07 +0800 CST
  • Martin Hope
    wimalopaan Somente operações bit a bit para std::byte em C++ 17? 2023-08-17 17:13:58 +0800 CST

Hot tag

python javascript c++ c# java typescript sql reactjs html

Explore

  • Início
  • Perguntas
    • Recentes
    • Highest score
  • tag
  • help

Footer

AskOverflow.Dev

About Us

  • About Us
  • Contact Us

Legal Stuff

  • Privacy Policy

Language

  • Pt
  • Server
  • Unix

© 2023 AskOverflow.DEV All Rights Reserve