SubmissionPublisher Klasse
Definition
Wichtig
Einige Informationen beziehen sich auf Vorabversionen, die vor dem Release ggf. grundlegend überarbeitet werden. Microsoft übernimmt hinsichtlich der hier bereitgestellten Informationen keine Gewährleistungen, seien sie ausdrücklich oder konkludent.
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
- 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 |
| 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 |
| 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 |
| SubmissionPublisher(IntPtr, JniHandleOwnership) |
Eine |
Eigenschaften
| Name | Beschreibung |
|---|---|
| Class |
Gibt die Laufzeitklasse dieses Werts |
| ClosedException |
Gibt die Ausnahme zurück |
| 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 |
| JniPeerMembers |
Eine |
| 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 |
| ThresholdClass |
Eine |
| ThresholdType |
Eine |
Methoden
| Name | Beschreibung |
|---|---|
| Clone() |
Erstellt und gibt eine Kopie dieses Objekts zurück. (Geerbt von Object) |
| Close() |
Sofern nicht bereits geschlossen, |
| CloseExceptionally(Throwable) |
Sofern nicht bereits geschlossen, |
| Construct(JniObjectReference, JniObjectReferenceOptions) |
Eine |
| 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 |
| Equals(Object) |
Eine |
| 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 |
| GetHashCode() |
Gibt einen Hashcodewert für das Objekt zurück. (Geerbt von Object) |
| IsSubscribed(Flow+ISubscriber) |
Eine |
| 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 |
| Offer(Object, Int64, TimeUnit, IBiPredicate) |
Veröffentlicht das angegebene Element nach Möglichkeit für jeden aktuellen Abonnent, indem er seine |
| SetHandle(IntPtr, JniHandleOwnership) |
Legt die Handle-Eigenschaft fest. (Geerbt von Object) |
| SetPeerReference(JniObjectReference, JniObjectReferenceOptions) |
Eine |
| Submit(Object) |
Veröffentlicht das angegebene Element für jeden aktuellen Abonnent, indem er seine |
| Subscribe(Flow+ISubscriber) |
Eine |
| 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 |
| IJavaPeerable.Finalized() |
Eine |
| IJavaPeerable.JniObjectReferenceControlBlock |
Eine |
| IJavaPeerable.SetJniIdentityHashCode(Int32) |
Eine |
| IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates) |
Eine |
| IJavaPeerable.SetPeerReference(JniObjectReference) |
Eine |
| IJavaPeerable.UnregisterFromRuntime() |
Eine |
Erweiterungsmethoden
| Name | Beschreibung |
|---|---|
| GetJniTypeName(IJavaPeerable) |
Ruft den JNI-Namen des Typs der Instanz |
| JavaAs<TResult>(IJavaPeerable) |
Versuchen Sie, die Eingabe |
| JavaCast<TResult>(IJavaObject) |
Führt eine android-laufzeitgecheckte Typkonvertierung aus. |
| JavaCast<TResult>(IJavaObject) |
Eine |
| TryJavaCast<TResult>(IJavaPeerable, TResult) |
Versuchen Sie, die Eingabe |