SubmissionPublisher Classe

Definição

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

[Android.Runtime.Register("java/util/concurrent/SubmissionPublisher", ApiSince=33, DoNotGenerateAcw=true)]
[Java.Interop.JavaTypeParameters(new System.String[] { "T" })]
public class SubmissionPublisher : Java.Lang.Object, IDisposable, Java.Interop.IJavaPeerable, Java.Lang.IAutoCloseable, Java.Util.Concurrent.Flow.IPublisher
[<Android.Runtime.Register("java/util/concurrent/SubmissionPublisher", ApiSince=33, DoNotGenerateAcw=true)>]
[<Java.Interop.JavaTypeParameters(new System.String[] { "T" })>]
type SubmissionPublisher = class
    inherit Object
    interface IAutoCloseable
    interface IJavaObject
    interface IDisposable
    interface IJavaPeerable
    interface Flow.IPublisher
Herança
SubmissionPublisher
Atributos
Implementações

Comentários

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado. Cada assinante atual recebe itens enviados recentemente na mesma ordem, a menos que sejam encontradas quedas ou exceções. O uso de um SubmissionPublisher permite que os geradores de itens atuem como fornecedores de fluxos reativos em conformidade que dependem do tratamento de drop e/ou do bloqueio para o controle de fluxo.

Um SubmissionPublisher usa o Executor fornecido em seu construtor para entrega aos assinantes. A melhor opção do Executor depende do uso esperado. Se os geradores de itens enviados forem executados em threads separados e o número de assinantes puder ser estimado, considere o uso de um Executors#newFixedThreadPool. Caso contrário, considere usar o padrão, normalmente o ForkJoinPool#commonPool.

O buffer permite que produtores e consumidores operem transitóriamente a taxas diferentes. Cada assinante usa um buffer independente. Os buffers são criados no primeiro uso e expandidos conforme necessário até o máximo fornecido. (A capacidade imposta pode ser arredondada até a potência mais próxima de dois e/ou limitada pelo maior valor suportado por essa implementação.) Invocações de não resultam diretamente na expansão do Flow.Subscription#request(long) request buffer, mas correm o risco de saturação se solicitações não preenchidas excederem a capacidade máxima. O valor padrão pode fornecer um ponto de Flow#defaultBufferSize() partida útil para escolher uma capacidade com base nas taxas, recursos e usos esperados.

Um único SubmissionPublisher pode ser compartilhado entre várias fontes. Ações em um thread de origem antes de publicar um item ou emitir um sinal <que eu>faço antes</i> ações subsequentes ao acesso correspondente por cada assinante. Mas estimativas relatadas de atraso e demanda foram projetadas para uso no monitoramento, não para controle de sincronização e podem refletir exibições obsoletas ou imprecisas do progresso.

Os métodos de publicação dão suporte a políticas diferentes sobre o que fazer quando os buffers estão saturados. O método #submit(Object) submit é bloqueado até que os recursos estejam disponíveis. Isso é mais simples, mas menos responsivo. Os offer métodos podem descartar itens (imediatamente ou com tempo limite limitado), mas oferecem uma oportunidade para interpor um manipulador e tentar novamente.

Se qualquer método assinante gerar uma exceção, sua assinatura será cancelada. Se um manipulador for fornecido como um argumento de construtor, ele será invocado antes do cancelamento em uma exceção no método Flow.Subscriber#onNext onNext, mas exceções em métodos Flow.Subscriber#onSubscribe onSubscribe, Flow.Subscriber#onError(Throwable) onError e Flow.Subscriber#onComplete() onComplete não será registrado ou tratado antes do cancelamento. Se o Executor fornecido gerar RejectedExecutionException (ou qualquer outro RuntimeException ou Error) ao tentar executar uma tarefa ou um manipulador de descarte gerar uma exceção ao processar um item descartado, a exceção será relançada. Nesses casos, nem todos os assinantes terão sido emitidos o item publicado. Geralmente, é uma boa prática para #closeExceptionally closeExceptionally esses casos.

O método #consume(Consumer) simplifica o suporte para um caso comum em que a única ação de um assinante é solicitar e processar todos os itens usando uma função fornecida.

Essa classe também pode servir como uma base conveniente para subclasses que geram itens e usar os métodos nessa classe para publicá-los. Por exemplo, aqui está uma classe que publica periodicamente os itens gerados de um fornecedor. (Na prática, você pode adicionar métodos para iniciar e parar a geração de forma independente, compartilhar Executores entre editores e assim por diante ou usar um SubmissionPublisher como um componente em vez de uma superclasse.)

{@code
            class PeriodicPublisher<T> extends SubmissionPublisher<T> {
              final ScheduledFuture<?> periodicTask;
              final ScheduledExecutorService scheduler;
              PeriodicPublisher(Executor executor, int maxBufferCapacity,
                                Supplier<? extends T> supplier,
                                long period, TimeUnit unit) {
                super(executor, maxBufferCapacity);
                scheduler = new ScheduledThreadPoolExecutor(1);
                periodicTask = scheduler.scheduleAtFixedRate(
                  () -> submit(supplier.get()), 0, period, unit);
              }
              public void close() {
                periodicTask.cancel(false);
                scheduler.shutdown();
                super.close();
              }
            }}

Aqui está um exemplo de uma Flow.Processor implementação. Ele usa solicitações de etapa única para seu editor para simplificar a ilustração. Uma versão mais adaptável poderia monitorar o fluxo usando a estimativa de retardo retornada, submitjuntamente com outros métodos utilitários.

{@code
            class TransformProcessor<S,T> extends SubmissionPublisher<T>
              implements Flow.Processor<S,T> {
              final Function<? super S, ? extends T> function;
              Flow.Subscription subscription;
              TransformProcessor(Executor executor, int maxBufferCapacity,
                                 Function<? super S, ? extends T> function) {
                super(executor, maxBufferCapacity);
                this.function = function;
              }
              public void onSubscribe(Flow.Subscription subscription) {
                (this.subscription = subscription).request(1);
              }
              public void onNext(S item) {
                subscription.request(1);
                submit(function.apply(item));
              }
              public void onError(Throwable ex) { closeExceptionally(ex); }
              public void onComplete() { close(); }
            }}

Adicionado em 9.

Java documentação para java.util.concurrent.SubmissionPublisher.

Partes desta página são modificações baseadas no trabalho criado e compartilhado pelo Project Open Source do Open Source e usadas de acordo com os termos descritos na Creative Commons 2.5.

Construtores

Nome Description
SubmissionPublisher()

Cria um novo SubmissionPublisher usando a ForkJoinPool#commonPool() entrega assíncrona para assinantes (a menos que não dê suporte a um nível de paralelismo de pelo menos dois, nesse caso, um novo Thread é criado para executar cada tarefa), com capacidade máxima de buffer e Flow#defaultBufferSizenenhum manipulador para exceções de Assinante no método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IExecutor, Int32, IBiConsumer)

Cria um novo SubmissionPublisher usando o Executor fornecido para entrega assíncrona aos assinantes, com o tamanho máximo de buffer fornecido para cada assinante e, se não nulo, o manipulador determinado invocado quando qualquer Assinante lança uma exceção no método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IExecutor, Int32)

Cria um novo SubmissionPublisher usando o Executor fornecido para entrega assíncrona aos assinantes, com o tamanho máximo de buffer fornecido para cada assinante e nenhum manipulador para exceções de Assinante no método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IntPtr, JniHandleOwnership)

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

Propriedades

Nome Description
Class

Retorna a classe de runtime deste Object.

(Herdado de Object)
ClosedException

Retorna a exceção associada #closeExceptionally(Throwable) closeExceptionallya , ou nulo, se não for fechada ou se fechada normalmente.

Executor

Retorna o Executor usado para entrega assíncrona.

Handle

O identificador para a instância subjacente do Android.

(Herdado de Object)
HasSubscribers

Retornará true se este editor tiver assinantes.

IsClosed

Retornará true se este editor não estiver aceitando envios.

JniIdentityHashCode

Obtém o código de hash de identidade atribuído a esse par Java pelo runtime de interoperabilidade.

(Herdado de Object)
JniPeerMembers

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

MaxBufferCapacity

Retorna a capacidade máxima de buffer por assinante.

NumberOfSubscribers

Retorna o número de assinantes atuais.

PeerReference

Obtém a referência de objeto JNI para este par Java.

(Herdado de Object)
Subscribers

Retorna uma lista de assinantes atuais para fins de monitoramento e acompanhamento, não para invocar Flow.Subscriber métodos nos assinantes.

ThresholdClass

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

ThresholdType

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

Métodos

Nome Description
Clone()

Cria e retorna uma cópia desse objeto.

(Herdado de Object)
Close()

A menos que já esteja fechado, emita Flow.Subscriber#onComplete() onComplete sinais para os assinantes atuais e não permite tentativas subsequentes de publicação.

CloseExceptionally(Throwable)

A menos que já esteja fechado, emita Flow.Subscriber#onError(Throwable) onError sinais para os assinantes atuais com o erro especificado e não permite tentativas subsequentes de publicação.

Consume(IConsumer)

Processa todos os itens publicados usando a função Consumidor fornecida.

Dispose()

Libera os recursos mantidos por esse par Java.

(Herdado de Object)
Dispose(Boolean)

Libera os recursos mantidos por esse par Java.

(Herdado de Object)
Equals(Object)

Indica se algum outro objeto é "igual a" este.

(Herdado de Object)
EstimateMaximumLag()

Retorna uma estimativa do número máximo de itens produzidos, mas ainda não consumidos entre todos os assinantes atuais.

EstimateMinimumDemand()

Retorna uma estimativa do número mínimo de itens solicitados (via Flow.Subscription#request(long) request) mas ainda não produzidos, entre todos os assinantes atuais.

GetHashCode()

Retorna um valor de código hash para o objeto.

(Herdado de Object)
IsSubscribed(Flow+ISubscriber)

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

JavaFinalize()

Chamado pelo coletor de lixo em um objeto quando a coleta de lixo determina que não há mais referências ao objeto.

(Herdado de Object)
Notify()

Ativa um único thread que está aguardando no monitor deste objeto.

(Herdado de Object)
NotifyAll()

Ativa todos os threads que estão aguardando no monitor deste objeto.

(Herdado de Object)
Offer(Object, IBiPredicate)

Publica o item determinado, se possível, para cada assinante atual invocando de forma assíncrona seu Flow.Subscriber#onNext(Object) onNext método.

Offer(Object, Int64, TimeUnit, IBiPredicate)

Publica o item determinado, se possível, para cada assinante atual invocando de forma assíncrona seu Flow.Subscriber#onNext(Object) onNext método, bloqueando enquanto os recursos de qualquer assinatura estão indisponíveis, até o tempo limite especificado ou até que o thread do chamador seja interrompido, momento em que o manipulador determinado (se não nulo) é invocado e, se ele retornar true, novamente repetido uma vez.

SetHandle(IntPtr, JniHandleOwnership)

Define a propriedade Handle

(Herdado de Object)
Submit(Object)

Publica o item especificado para cada assinante atual invocando de forma assíncrona seu Flow.Subscriber#onNext(Object) onNext método, bloqueando de forma ininterrupta enquanto os recursos para qualquer assinante não estão disponíveis.

Subscribe(Flow+ISubscriber)

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

ToArray<T>()

Cria uma matriz gerenciada com base nesse wrapper de matriz Java.

(Herdado de Object)
ToString()

Retorna uma representação de cadeia de caracteres do objeto.

(Herdado de Object)
UnregisterFromRuntime()

Cancela o registro desse par Java do runtime de interoperabilidade.

(Herdado de Object)
Wait()

Faz com que o thread atual aguarde até ser despertado, normalmente por ser <notificado/em> ou <em>interrompido</em>.<>

(Herdado de Object)
Wait(Int64, Int32)

Faz com que o thread atual aguarde até que ele seja despertado, normalmente por ser <>notificado</em> ou <em>interrompido</em>, ou até que uma determinada quantidade de tempo real tenha decorrido.

(Herdado de Object)
Wait(Int64)

Faz com que o thread atual aguarde até que ele seja despertado, normalmente por ser <>notificado</em> ou <em>interrompido</em>, ou até que uma determinada quantidade de tempo real tenha decorrido.

(Herdado de Object)

Implantações explícitas de interface

Nome Description
IJavaPeerable.Disposed()

Notifica o runtime de interoperabilidade de que esse par gerenciado foi descartado.

(Herdado de Object)
IJavaPeerable.DisposeUnlessReferenced()

Libera esse par Java, a menos que seja mantido por outra referência gerenciada.

(Herdado de Object)
IJavaPeerable.Finalized()

Notifica o runtime de interoperabilidade de que esse par gerenciado foi finalizado.

(Herdado de Object)
IJavaPeerable.JniManagedPeerState

Obtém o estado que descreve a relação entre esse par gerenciado e sua referência de JNI.

(Herdado de Object)
IJavaPeerable.SetJniIdentityHashCode(Int32)

Define o código de hash de identidade de interoperabilidade para esse par Java.

(Herdado de Object)
IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates)

Define o estado de par gerenciado e JNI para esse par Java.

(Herdado de Object)
IJavaPeerable.SetPeerReference(JniObjectReference)

Define a referência de objeto JNI usada por esse par gerenciado.

(Herdado de Object)

Métodos de Extensão

Nome Description
GetJniTypeName(IJavaPeerable)

Obtém o nome JNI do tipo da instância self.

JavaAs<TResult>(IJavaPeerable)

Tente coagir a digitar selfTResult, verificando se a coerção é válida no lado Java.

JavaCast<TResult>(IJavaObject)

Executa uma conversão de tipo marcada por runtime do Android.

JavaCast<TResult>(IJavaObject)

Um Flow.Publisher que emite de forma assíncrona itens enviados (não nulos) aos assinantes atuais até que seja fechado.

TryJavaCast<TResult>(IJavaPeerable, TResult)

Tente coagir a digitar selfTResult, verificando se a coerção é válida no lado Java.

Aplica-se a