SubmissionPublisher Classe
Definição
Importante
Algumas informações se referem a produtos de pré-lançamento que podem ser substancialmente modificados antes do lançamento. A Microsoft não oferece garantias, expressas ou implícitas, das informações aqui fornecidas.
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
- 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 |
| 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 |
| 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 |
| SubmissionPublisher(IntPtr, JniHandleOwnership) |
Um |
Propriedades
| Nome | Description |
|---|---|
| Class |
Retorna a classe de runtime deste |
| ClosedException |
Retorna a exceção associada |
| 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 |
| 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 |
| ThresholdClass |
Um |
| ThresholdType |
Um |
Métodos
| Nome | Description |
|---|---|
| Clone() |
Cria e retorna uma cópia desse objeto. (Herdado de Object) |
| Close() |
A menos que já esteja fechado, emita |
| CloseExceptionally(Throwable) |
A menos que já esteja fechado, emita |
| 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 |
| GetHashCode() |
Retorna um valor de código hash para o objeto. (Herdado de Object) |
| IsSubscribed(Flow+ISubscriber) |
Um |
| 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 |
| Offer(Object, Int64, TimeUnit, IBiPredicate) |
Publica o item determinado, se possível, para cada assinante atual invocando de forma assíncrona seu |
| 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 |
| Subscribe(Flow+ISubscriber) |
Um |
| 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 |
| JavaAs<TResult>(IJavaPeerable) |
Tente coagir a digitar |
| JavaCast<TResult>(IJavaObject) |
Executa uma conversão de tipo marcada por runtime do Android. |
| JavaCast<TResult>(IJavaObject) |
Um |
| TryJavaCast<TResult>(IJavaPeerable, TResult) |
Tente coagir a digitar |