Sprache

SubmissionPublisher Klasse

Definition

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

[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
Vererbung
SubmissionPublisher
Attribute
Implementiert

Hinweise

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist. Jeder aktuelle Abonnent empfängt neu übermittelte Elemente in derselben Reihenfolge, es sei denn, Es treten Tropfen oder Ausnahmen auf. Mit einem SubmissionPublisher können Elementgeneratoren als kompatible reaktive Datenströme fungieren, die auf drop handling und/oder blockieren für die Flusssteuerung vertrauen.

Ein SubmissionPublisher verwendet den Executor im Konstruktor bereitgestellten Für die Übermittlung an Abonnenten. Die beste Wahl des Executors hängt von der erwarteten Verwendung ab. Wenn die Generatoren von übermittelten Elementen in separaten Threads ausgeführt werden und die Anzahl der Abonnenten geschätzt werden kann, erwägen Sie die Verwendung einer Executors#newFixedThreadPool. Ziehen Sie andernfalls die Verwendung der Standardeinstellung in Betracht, normalerweise die ForkJoinPool#commonPool.

Die Pufferung ermöglicht es Produzenten und Verbrauchern, vorübergehend mit unterschiedlichen Tarifen zu arbeiten. Jeder Abonnent verwendet einen unabhängigen Puffer. Puffer werden bei der ersten Verwendung erstellt und nach Bedarf bis zum angegebenen Maximum erweitert. (Die erzwungene Kapazität kann auf die nächste Leistung von zwei und/oder an den größten Von dieser Implementierung unterstützten Wert aufgerundet werden.) Aufrufe führen Flow.Subscription#request(long) request nicht direkt zur Puffererweiterung, aber risikosättigung, wenn nicht ausgefüllte Anforderungen die maximale Kapazität überschreiten. Der Standardwert kann Flow#defaultBufferSize() einen nützlichen Ausgangspunkt für die Auswahl einer Kapazität basierend auf erwarteten Raten, Ressourcen und Nutzungen bieten.

Ein einzelner SubmissionPublisher kann für mehrere Quellen freigegeben werden. Aktionen in einem Quellthread vor der Veröffentlichung eines Elements oder Ausstellen eines Signals<, das vor></i erfolgt,/i> Aktionen, die auf den entsprechenden Zugriff durch jeden Abonnenten folgt. Gemeldete Schätzungen von Verzögerung und Nachfrage sind jedoch für die Verwendung in der Überwachung vorgesehen, nicht für die Synchronisierungskontrolle und können veraltete oder ungenaue Ansichten des Fortschritts widerspiegeln.

Publikationsmethoden unterstützen unterschiedliche Richtlinien dazu, was zu tun ist, wenn Puffer gesättigt sind. Die Methode #submit(Object) submit blockiert, bis Ressourcen verfügbar sind. Dies ist am einfachsten, aber am wenigsten reaktionsfähig. Die offer Methoden können Elemente ablegen (entweder sofort oder mit gebundenem Timeout), bieten aber die Möglichkeit, einen Handler zu interposieren und dann erneut zu versuchen.

Wenn eine Subscriber-Methode eine Ausnahme auslöst, wird das Abonnement gekündigt. Wenn ein Handler als Konstruktorargument bereitgestellt wird, wird er vor dem Abbruch bei einer Ausnahme in der Methode Flow.Subscriber#onNext onNextaufgerufen, aber Ausnahmen in Methoden Flow.Subscriber#onSubscribe onSubscribeund Flow.Subscriber#onError(Throwable) onErrorFlow.Subscriber#onComplete() onComplete werden vor dem Abbruch nicht aufgezeichnet oder behandelt. Wenn der angegebene Executor beim Versuch, eine Aufgabe auszuführen, eine Ausnahme auslöst RejectedExecutionException (oder ein anderer RuntimeException oder Error), oder ein Drop-Handler löst beim Verarbeiten eines verworfenen Elements eine Ausnahme aus, dann wird die Ausnahme erneut ausgelöst. In diesen Fällen wurden nicht alle Abonnenten das veröffentlichte Element ausgestellt. In diesen Fällen ist es in der Regel empfehlenswert #closeExceptionally closeExceptionally .

Die Methode #consume(Consumer) vereinfacht die Unterstützung für einen gängigen Fall, in dem die einzige Aktion eines Abonnenten darin besteht, alle Elemente mithilfe einer bereitgestellten Funktion anzufordern und zu verarbeiten.

Diese Klasse kann auch als praktische Basis für Unterklassen dienen, die Elemente generieren, und die Methoden in dieser Klasse verwenden, um sie zu veröffentlichen. Hier ist beispielsweise eine Klasse, die die von einem Lieferanten generierten Elemente regelmäßig veröffentlicht. (In der Praxis fügen Sie Methoden zum unabhängigen Starten und Beenden der Generierung hinzu, um Executors unter Herausgebern usw. freizugeben oder einen SubmissionPublisher anstelle einer Superklasse als Komponente zu verwenden.)

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

Hier ist ein Beispiel für eine Flow.Processor Implementierung. Er verwendet Einzelschrittanforderungen an seinen Herausgeber zur Vereinfachung der Illustration. Eine adaptivere Version könnte den Fluss überwachen, indem die von submitder Verzögerung zurückgegebene Schätzung zusammen mit anderen Hilfsmethoden verwendet wird.

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

Hinzugefügt in 9.

Java Dokumentation für java.util.concurrent.SubmissionPublisher.

Teile dieser Seite sind Änderungen auf der Grundlage von Arbeiten, die von der Android Open Source Project erstellt und gemeinsam verwendet und gemäß den in der 2.5 Attribution License beschriebenen Begriffen verwendet werden.

Konstruktoren

Name Beschreibung
SubmissionPublisher()

Erstellt einen neuen SubmissionPublisher unter Verwendung der ForkJoinPool#commonPool() asynchronen Übermittlung an Abonnenten (es sei denn, es unterstützt keine Parallelitätsstufe von mindestens zwei, in diesem Fall wird ein neuer Thread erstellt, um jede Aufgabe auszuführen), mit maximaler Pufferkapazität von Flow#defaultBufferSize, und kein Handler für Abonnentenausnahmen in der Methode Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IExecutor, Int32, IBiConsumer)

Erstellt einen neuen SubmissionPublisher mit dem angegebenen Executor für die asynchrone Übermittlung an Abonnenten, mit der angegebenen maximalen Puffergröße für jeden Abonnent, und wenn der angegebene Handler nicht null ist, wird der angegebene Handler aufgerufen, wenn ein Subscriber eine Ausnahme in der Methode Flow.Subscriber#onNext(Object) onNextauslöst.

SubmissionPublisher(IExecutor, Int32)

Erstellt einen neuen SubmissionPublisher mit dem angegebenen Executor für die asynchrone Übermittlung an Abonnenten mit der angegebenen maximalen Puffergröße für jeden Abonnenten und keinen Handler für Abonnentenausnahmen in der Methode Flow.Subscriber#onNext(Object) onNext.

SubmissionPublisher(IntPtr, JniHandleOwnership)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

Eigenschaften

Name Beschreibung
Class

Gibt die Laufzeitklasse dieses Werts Objectzurück.

(Geerbt von Object)
ClosedException

Gibt die Ausnahme zurück #closeExceptionally(Throwable) closeExceptionally, die mit oder null verknüpft ist, wenn sie nicht geschlossen oder normal geschlossen ist.

Executor

Gibt den Executor zurück, der für die asynchrone Übermittlung verwendet wird.

Handle

Das Handle für die zugrunde liegende Android-Instanz.

(Geerbt von Object)
HasSubscribers

Gibt true zurück, wenn dieser Herausgeber abonnenten hat.

IsClosed

Gibt true zurück, wenn dieser Herausgeber keine Übermittlungen akzeptiert.

JniIdentityHashCode

Ruft den Identitätshashcode ab, der diesem Java Peer von der Interop-Laufzeit zugewiesen ist.

(Geerbt von Object)
JniManagedPeerState

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
JniPeerMembers

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

MaxBufferCapacity

Gibt die maximale Pufferkapazität pro Abonnent zurück.

NumberOfSubscribers

Gibt die Anzahl der aktuellen Abonnenten zurück.

PeerReference

Ruft den JNI-Objektverweis für diesen Java Peer ab.

(Geerbt von Object)
Subscribers

Gibt eine Liste der aktuellen Abonnenten für Überwachungs- und Nachverfolgungszwecke zurück, nicht zum Aufrufen von Flow.Subscriber Methoden für die Abonnenten.

ThresholdClass

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

ThresholdType

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

Methoden

Name Beschreibung
Clone()

Erstellt und gibt eine Kopie dieses Objekts zurück.

(Geerbt von Object)
Close()

Sofern nicht bereits geschlossen, Flow.Subscriber#onComplete() onComplete gibt es Signale an aktuelle Abonnenten und verbietet nachfolgende Veröffentlichungsversuche.

CloseExceptionally(Throwable)

Sofern nicht bereits geschlossen, Flow.Subscriber#onError(Throwable) onError gibt es Signale an aktuelle Abonnenten mit dem gegebenen Fehler und verbietet nachfolgende Veröffentlichungsversuche.

Construct(JniObjectReference, JniObjectReferenceOptions)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
Consume(IConsumer)

Verarbeitet alle veröffentlichten Elemente mithilfe der angegebenen Consumer-Funktion.

Dispose()

Veröffentlicht die Ressourcen, die von diesem Java Peer gehalten werden.

(Geerbt von Object)
Dispose(Boolean)

Veröffentlicht die Ressourcen, die von diesem Java Peer gehalten werden.

(Geerbt von Object)
DisposeUnlessReferenced()

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
Equals(Object)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
Equals(Object)

Gibt an, ob ein anderes Objekt "gleich" diesem Objekt ist.

(Geerbt von Object)
EstimateMaximumLag()

Gibt eine Schätzung der maximalen Anzahl der produzierten, aber noch nicht verbrauchten Elemente unter allen aktuellen Abonnenten zurück.

EstimateMinimumDemand()

Gibt eine Schätzung der Mindestanzahl der angeforderten Elemente (via Flow.Subscription#request(long) request) zurück, aber noch nicht produziert, unter allen aktuellen Abonnenten.

GetHashCode()

Gibt einen Hashcodewert für das Objekt zurück.

(Geerbt von Object)
IsSubscribed(Flow+ISubscriber)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

JavaFinalize()

Wird vom Garbage Collector für ein Objekt aufgerufen, wenn die Garbage Collection bestimmt, dass keine weiteren Verweise auf das Objekt vorhanden sind.

(Geerbt von Object)
Notify()

Aktiviert einen einzelnen Thread, der auf dem Monitor dieses Objekts wartet.

(Geerbt von Object)
NotifyAll()

Aktiviert alle Threads, die auf dem Monitor dieses Objekts warten.

(Geerbt von Object)
Offer(Object, IBiPredicate)

Veröffentlicht das angegebene Element, falls möglich, für jeden aktuellen Abonnent, indem er seine Flow.Subscriber#onNext(Object) onNext Methode asynchron aufruft.

Offer(Object, Int64, TimeUnit, IBiPredicate)

Veröffentlicht das angegebene Element nach Möglichkeit für jeden aktuellen Abonnent, indem er seine Flow.Subscriber#onNext(Object) onNext Methode asynchron aufruft, sperrt, während Ressourcen für ein Abonnement nicht verfügbar sind, bis zum angegebenen Timeout oder bis der Aufruferthread unterbrochen wird, an welchem Punkt der angegebene Handler (wenn nicht null) aufgerufen wird, und wenn er "true" zurückgibt, einmal wiederholt.

SetHandle(IntPtr, JniHandleOwnership)

Legt die Handle-Eigenschaft fest.

(Geerbt von Object)
SetPeerReference(JniObjectReference, JniObjectReferenceOptions)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
Submit(Object)

Veröffentlicht das angegebene Element für jeden aktuellen Abonnent, indem er seine Flow.Subscriber#onNext(Object) onNext Methode asynchron aufruft und die Ressourcen für jeden Abonnenten nicht verfügbar ist.

Subscribe(Flow+ISubscriber)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

ToArray<T>()

Erstellt ein verwaltetes Array aus diesem Java Arraywrapper.

(Geerbt von Object)
ToString()

Gibt eine Zeichenfolgendarstellung des Objekts zurück.

(Geerbt von Object)
UnregisterFromRuntime()

Hebt die Registrierung dieses Java Peers aus der Interop-Laufzeit auf.

(Geerbt von Object)
Wait()

Bewirkt, dass der aktuelle Thread wartet, bis er wach ist, in der Regel durch em benachrichtigt/em< oder >em<unterbrochen>/em<.><>

(Geerbt von Object)
Wait(Int64, Int32)

Bewirkt, dass der aktuelle Thread wartet, bis er wach ist, in der Regel durch <em>benachrichtigt</em> oder <em>unterbrochen</em> oder bis eine bestimmte Menge an Echtzeit verstrichen ist.

(Geerbt von Object)
Wait(Int64)

Bewirkt, dass der aktuelle Thread wartet, bis er wach ist, in der Regel durch <em>benachrichtigt</em> oder <em>unterbrochen</em> oder bis eine bestimmte Menge an Echtzeit verstrichen ist.

(Geerbt von Object)

Explizite Schnittstellenimplementierungen

Name Beschreibung
IJavaPeerable.Disposed()

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.Finalized()

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.JniObjectReferenceControlBlock

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.SetJniIdentityHashCode(Int32)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.SetPeerReference(JniObjectReference)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

(Geerbt von JavaObject)
IJavaPeerable.UnregisterFromRuntime()

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

Erweiterungsmethoden

Name Beschreibung
GetJniTypeName(IJavaPeerable)

Ruft den JNI-Namen des Typs der Instanz selfab.

JavaAs<TResult>(IJavaPeerable)

Versuchen Sie, die Eingabe selfzu TResult erzwingen, und überprüfen Sie, ob die Koersion auf der Java Seite gültig ist.

JavaCast<TResult>(IJavaObject)

Führt eine android-laufzeitgecheckte Typkonvertierung aus.

JavaCast<TResult>(IJavaObject)

Eine Flow.Publisher , die übermittelte (nicht NULL)-Elemente asynchron an aktuelle Abonnenten ausgibt, bis sie geschlossen ist.

TryJavaCast<TResult>(IJavaPeerable, TResult)

Versuchen Sie, die Eingabe selfzu TResult erzwingen, und überprüfen Sie, ob die Koersion auf der Java Seite gültig ist.

Gilt für: