Lenguaje

SubmissionPublisher Clase

Definición

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

[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.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
Herencia
SubmissionPublisher
Atributos
Implementaciones

Comentarios

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra. Cada suscriptor actual recibe los elementos recién enviados en el mismo orden a menos que se encuentren caídas o excepciones. El uso de submissionPublisher permite que los generadores de elementos actúen como publicadores reactivos compatibles que dependen del control de caídas o del bloqueo para el control de flujo.

Un objeto SubmissionPublisher usa el Executor proporcionado en su constructor para su entrega a los suscriptores. La mejor opción de Executor depende del uso esperado. Si los generadores de elementos enviados se ejecutan en subprocesos independientes y se puede estimar el número de suscriptores, considere la posibilidad de usar .Executors#newFixedThreadPool De lo contrario, considere la posibilidad de usar el valor predeterminado, normalmente .ForkJoinPool#commonPool

El almacenamiento en búfer permite a los productores y consumidores operar transitoriamente a diferentes velocidades. Cada suscriptor usa un búfer independiente. Los búferes se crean al usar por primera vez y se expanden según sea necesario hasta el máximo especificado. (La capacidad aplicada se puede redondear hasta la potencia más cercana de dos o limitadas por el valor más grande admitido por esta implementación). Las invocaciones de Flow.Subscription#request(long) request no generan directamente la expansión del búfer, pero la saturación del riesgo si las solicitudes no rellenadas superan la capacidad máxima. El valor predeterminado de puede proporcionar un punto de Flow#defaultBufferSize() partida útil para elegir una capacidad en función de las tasas, los recursos y los usos esperados.

Un único SubmissionPublisher puede compartirse entre varios orígenes. Acciones en un subproceso de origen antes de publicar un elemento o emitir una señal <i>acciones< anteriores/i> posteriores al acceso correspondiente por cada suscriptor. Pero las estimaciones notificadas de retraso y demanda están diseñadas para su uso en la supervisión, no para el control de sincronización, y pueden reflejar vistas obsoletas o inexactas del progreso.

Los métodos de publicación admiten diferentes directivas sobre qué hacer cuando los búferes están saturados. Los bloques de métodos #submit(Object) submit se bloquean hasta que los recursos estén disponibles. Esto es más sencillo, pero menos sensible. Los offer métodos pueden quitar elementos (inmediatamente o con tiempo de espera limitado), pero proporcionan una oportunidad para interponer un controlador y, a continuación, volver a intentarlo.

Si algún método de suscriptor produce una excepción, se cancela su suscripción. Si se proporciona un controlador como argumento de constructor, se invoca antes de la cancelación tras una excepción en el método Flow.Subscriber#onNext onNext, pero las excepciones de los métodos Flow.Subscriber#onSubscribe onSubscribe, Flow.Subscriber#onError(Throwable) onError y Flow.Subscriber#onComplete() onComplete no se registran ni controlan antes de la cancelación. Si el ejecutor proporcionado inicia RejectedExecutionException (o cualquier otro RuntimeException o Error) al intentar ejecutar una tarea, o un controlador de colocación produce una excepción al procesar un elemento descartado, la excepción se vuelve a iniciar. En estos casos, no todos los suscriptores habrán sido emitidos el elemento publicado. Normalmente, es recomendable hacerlo #closeExceptionally closeExceptionally en estos casos.

El método #consume(Consumer) simplifica la compatibilidad con un caso común en el que la única acción de un suscriptor es solicitar y procesar todos los elementos mediante una función proporcionada.

Esta clase también puede servir como una base cómoda para las subclases que generan elementos y usar los métodos de esta clase para publicarlos. Por ejemplo, esta es una clase que publica periódicamente los elementos generados a partir de un proveedor. (En la práctica, puede agregar métodos para iniciar y detener la generación de forma independiente, para compartir ejecutores entre publicadores, etc., o usar un SubmissionPublisher como componente en lugar de una superclase).

{@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();
              }
            }}

Este es un ejemplo de una Flow.Processor implementación. Usa solicitudes de un solo paso a su publicador para simplificar la ilustración. Una versión más adaptable podría supervisar el flujo mediante la estimación de retraso devuelta de submit, junto con otros métodos de utilidad.

{@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(); }
            }}

Agregado en 9.

Java documentación para java.util.concurrent.SubmissionPublisher.

Las partes de esta página son modificaciones basadas en el trabajo creado y compartido por el Android y se usan según los términos descritos en creative Creative Commons 2.5 Attribution License.

Constructores

Nombre Description
SubmissionPublisher()

Crea un nuevo objeto SubmissionPublisher mediante para ForkJoinPool#commonPool() la entrega asincrónica a los suscriptores (a menos que no admita un nivel de paralelismo de al menos dos, en cuyo caso se crea un nuevo subproceso para ejecutar cada tarea), con capacidad máxima de búfer de Flow#defaultBufferSizey sin controlador para excepciones de suscriptor en el método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IExecutor, Int32, IBiConsumer)

Crea un nuevo objeto SubmissionPublisher mediante el ejecutor especificado para la entrega asincrónica a los suscriptores, con el tamaño máximo de búfer especificado para cada suscriptor y, si no es NULL, el controlador especificado invoca cuando cualquier suscriptor produce una excepción en el método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IExecutor, Int32)

Crea un nuevo objeto SubmissionPublisher mediante el ejecutor especificado para la entrega asincrónica a los suscriptores, con el tamaño máximo de búfer especificado para cada suscriptor y ningún controlador para las excepciones de suscriptor en el método Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IntPtr, JniHandleOwnership)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

Propiedades

Nombre Description
Class

Devuelve la clase en tiempo de ejecución de este Objectobjeto .

(Heredado de Object)
ClosedException

Devuelve la excepción asociada a #closeExceptionally(Throwable) closeExceptionally, o null si no se cierra o si se cierra normalmente.

Executor

Devuelve el ejecutor usado para la entrega asincrónica.

Handle

Identificador de la instancia de Android subyacente.

(Heredado de Object)
HasSubscribers

Devuelve true si este publicador tiene suscriptores.

IsClosed

Devuelve true si este publicador no acepta envíos.

JniIdentityHashCode

Obtiene el código hash de identidad asignado a este Java del mismo nivel por el tiempo de ejecución de interoperabilidad.

(Heredado de Object)
JniManagedPeerState

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
JniPeerMembers

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

MaxBufferCapacity

Devuelve la capacidad máxima del búfer por suscriptor.

NumberOfSubscribers

Devuelve el número de suscriptores actuales.

PeerReference

Obtiene la referencia de objeto JNI para este Java del mismo nivel.

(Heredado de Object)
Subscribers

Devuelve una lista de suscriptores actuales con fines de supervisión y seguimiento, no para invocar Flow.Subscriber métodos en los suscriptores.

ThresholdClass

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

ThresholdType

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

Métodos

Nombre Description
Clone()

Crea y devuelve una copia de este objeto.

(Heredado de Object)
Close()

A menos que ya esté cerrado, emite Flow.Subscriber#onComplete() onComplete señales a los suscriptores actuales y no permite que los intentos posteriores se publiquen.

CloseExceptionally(Throwable)

A menos que ya se haya cerrado, emite Flow.Subscriber#onError(Throwable) onError señales a los suscriptores actuales con el error dado y no permite que los intentos posteriores se publiquen.

Construct(JniObjectReference, JniObjectReferenceOptions)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
Consume(IConsumer)

Procesa todos los elementos publicados mediante la función Consumidor determinada.

Dispose()

Libera los recursos mantenidos por este Java del mismo nivel.

(Heredado de Object)
Dispose(Boolean)

Libera los recursos mantenidos por este Java del mismo nivel.

(Heredado de Object)
DisposeUnlessReferenced()

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
Equals(Object)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
Equals(Object)

Indica si algún otro objeto es "igual a" este.

(Heredado de Object)
EstimateMaximumLag()

Devuelve una estimación del número máximo de elementos producidos pero aún no consumidos entre todos los suscriptores actuales.

EstimateMinimumDemand()

Devuelve una estimación del número mínimo de elementos solicitados (a través Flow.Subscription#request(long) requestde ) pero aún no generados, entre todos los suscriptores actuales.

GetHashCode()

Devuelve un valor de código hash del objeto.

(Heredado de Object)
IsSubscribed(Flow+ISubscriber)

Devuelve true si el suscriptor especificado está suscrito actualmente.

JavaFinalize()
Obsoletos.

Lo llama el recolector de elementos no utilizados en un objeto cuando la recolección de elementos no utilizados determina que no hay más referencias al objeto .

(Heredado de Object)
Notify()

Activa un único subproceso que está esperando en el monitor de este objeto.

(Heredado de Object)
NotifyAll()

Activa todos los subprocesos que están esperando en el monitor de este objeto.

(Heredado de Object)
Offer(Object, IBiPredicate)

Publica el elemento especificado, si es posible, en cada suscriptor actual invocando de forma asincrónica su Flow.Subscriber#onNext(Object) onNext método.

Offer(Object, Int64, TimeUnit, IBiPredicate)

Publica el elemento especificado, si es posible, para cada suscriptor actual invocando de forma asincrónica su Flow.Subscriber#onNext(Object) onNext método, bloqueando mientras los recursos de cualquier suscripción no están disponibles, hasta el tiempo de espera especificado o hasta que se interrumpa el subproceso del autor de la llamada, en cuyo punto se invoca el controlador especificado (si no es null) y, si devuelve true, se reintenta una vez.

SetHandle(IntPtr, JniHandleOwnership)

Establece la propiedad Handle.

(Heredado de Object)
SetPeerReference(JniObjectReference, JniObjectReferenceOptions)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
Submit(Object)

Publica el elemento especificado en cada suscriptor actual invocando de forma asincrónica su Flow.Subscriber#onNext(Object) onNext método, bloqueando de forma ininterrumpida mientras los recursos de cualquier suscriptor no están disponibles.

Subscribe(Flow+ISubscriber)

Agrega el suscriptor especificado a menos que ya esté suscrito.

ToArray<T>()

Crea una matriz administrada a partir de este contenedor de matriz Java.

(Heredado de Object)
ToString()

Devuelve una representación de cadena del objeto.

(Heredado de Object)
UnregisterFromRuntime()

Anula el registro de este Java del mismo nivel del tiempo de ejecución de interoperabilidad.

(Heredado de Object)
Wait()

Hace que el subproceso actual espere hasta que se despierta, normalmente por ser em notificado/em< o >em<interrumpido>/em<.><>

(Heredado de Object)
Wait(Int64, Int32)

Hace que el subproceso actual espere hasta que se despierte, normalmente por ser <em>notificado</em> o <em>interrumpido</em>, o hasta que haya transcurrido una cierta cantidad de tiempo real.

(Heredado de Object)
Wait(Int64)

Hace que el subproceso actual espere hasta que se despierte, normalmente por ser <em>notificado</em> o <em>interrumpido</em>, o hasta que haya transcurrido una cierta cantidad de tiempo real.

(Heredado de Object)

Implementaciones de interfaz explícitas

Nombre Description
IJavaPeerable.Disposed()

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.Finalized()

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.JniObjectReferenceControlBlock

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.SetJniIdentityHashCode(Int32)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.SetPeerReference(JniObjectReference)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

(Heredado de JavaObject)
IJavaPeerable.UnregisterFromRuntime()

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

Métodos de extensión

Nombre Description
GetJniTypeName(IJavaPeerable)

Obtiene el nombre JNI del tipo de la instancia self.

JavaAs<TResult>(IJavaPeerable)

Intente coerción self para escribir TResult, comprobando que la coerción es válida en el lado de Java.

JavaCast<TResult>(IJavaObject)

Realiza una conversión de tipos comprobados en tiempo de ejecución de Android.

JavaCast<TResult>(IJavaObject)

que Flow.Publisher emite de forma asincrónica elementos enviados (no NULL) a los suscriptores actuales hasta que se cierra.

TryJavaCast<TResult>(IJavaPeerable, TResult)

Intente coerción self para escribir TResult, comprobando que la coerción es válida en el lado de Java.

Se aplica a