語言

SubmissionPublisher 類別

定義

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

[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
繼承
SubmissionPublisher
屬性
實作

備註

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。 每位現有訂閱者都會以相同順序收到新提交的商品,除非遇到掉落或例外情況。 使用 SubmissionPublisher 讓項目產生器能作為合規 的反應式串流 ,依賴 drop 處理和/或阻擋來控制流量。

SubmissionPublisher 利用其建構器中提供的資訊 Executor 來傳遞給訂閱者。 執行者的最佳選擇取決於預期使用情況。 如果提交項目的產生器在獨立執行緒中運行,且訂閱者數量可以估計,請考慮使用 。Executors#newFixedThreadPool 否則請考慮使用預設值,通常是 ForkJoinPool#commonPool。

緩衝讓生產者與消費者能暫時以不同速率運作。 每個用戶使用獨立的緩衝區。 緩衝區會在首次使用時建立,並根據需要擴展至指定最大值。 (強制容量可向上取整至最接近的二的冪次方,和/或以此實作支持的最大值為界。)調用 的 Flow.Subscription#request(long) request 呼叫不會直接導致緩衝區擴充,但若未填滿的請求超過最大容量,則有飽和的風險。 預設值 可 Flow#defaultBufferSize() 作為根據預期速率、資源及使用量選擇容量的有用起點。

單一 SubmissionPublisher 可能在多個來源間共享。 在發布項目或發出 i<happen-before>/i< 訊號>前,來源執行緒中每個訂閱者在相應存取後的動作。 但報告的延遲與需求估計是用於監控,而非同步控制,且可能反映陳舊或不準確的進展觀點。

發表方法支持不同政策,說明緩衝區飽和時應採取的措施。 方法 #submit(Object) submit 區塊直到資源可用。 這是最簡單,但反應最慢的。 這些 offer 方法可能會立即丟棄項目(或有界逾時),但提供介入處理程序並重試的機會。

若任何訂閱者方法拋出例外,該訂閱即被取消。 若以建構子參數形式提供處理器,則在方法中異常Flow.Subscriber#onNext onNext時,該處理器會在取消前被呼叫,但方法Flow.Subscriber#onSubscribe onSubscribeFlow.Subscriber#onError(Throwable) onErrorFlow.Subscriber#onComplete() onComplete中的異常在刪除前不會被記錄或處理。 若所提供的執行器在執行任務時拋 RejectedExecutionException 出(或其他任何 RuntimeException 或錯誤),或丟棄處理程序在處理丟棄項目時拋出異常,則該例外會被重新拋出。 在這些情況下,並非所有訂閱者都已獲得已發佈的項目。 在這種情況下,通常 #closeExceptionally closeExceptionally 這樣做是很好的做法。

此方法 #consume(Consumer) 簡化了一種常見情境的支援,即訂閱者唯一的操作是使用所提供的函式請求並處理所有項目。

此類別亦可作為子類別產生項目的便利基礎,並利用該類別的方法發佈。 舉例來說,這裡有一個類別會定期發布供應商產生的項目。 (實務上,你可以加入獨立啟動與停止生成的方法,在出版者間共享執行者等,或將 SubmissionPublisher 作為元件而非超類別。)

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

這裡有一個實作範例 Flow.Processor 。 它採用單步請求給出版商,以簡化說明。 較適應性的版本則可利用從 submit回傳的延遲估計值監控流量,並結合其他實用性方法。

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

新增9。

Java 文件 java.util.concurrent.SubmissionPublisher。

本頁部分內容為基於 Open Source Project 所創建與分享的作品,並依授權條款所描述的使用進行修改。

建構函式

名稱 Description
SubmissionPublisher()

使用 for ForkJoinPool#commonPool() async 傳送給訂閱者(除非不支援至少二級平行處理,則需建立新執行緒執行每個任務),最大緩衝容量為 Flow#defaultBufferSize,且方法中不設訂閱者例外 Flow.Subscriber#onNext(Object) onNext處理程序。

SubmissionPublisher(IExecutor, Int32, IBiConsumer)

使用給定的執行器建立新的 SubmissionPublisher,進行非同步傳遞給訂閱者,每個訂閱者的最大緩衝區大小為指定,若非空,則當任何訂閱者在方法 Flow.Subscriber#onNext(Object) onNext中拋出例外時,會呼叫該處理器。

SubmissionPublisher(IExecutor, Int32)

使用給定的執行器建立新的 SubmissionPublisher,進行非同步傳遞給訂閱者,每個訂閱者擁有最大緩衝區大小,且方法中不包含訂閱者例外 Flow.Subscriber#onNext(Object) onNext的處理器。

SubmissionPublisher(IntPtr, JniHandleOwnership)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

屬性

名稱 Description
Class

傳回這個 Object的運行時間類別。

(繼承來源 Object)
ClosedException

回傳與 #closeExceptionally(Throwable) closeExceptionally相關聯的例外,若未封閉或正常閉則為 null。

Executor

回傳用於非同步傳送的執行者。

Handle

基礎Android實例的句柄。

(繼承來源 Object)
HasSubscribers

如果該出版商有任何訂閱者,則回傳為真。

IsClosed

如果該出版社不接受投稿,則會回傳為真。

JniIdentityHashCode

取得由互通執行時指派給此 Java 對等端的身份雜湊碼。

(繼承來源 Object)
JniManagedPeerState

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
JniPeerMembers

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

MaxBufferCapacity

回傳每位用戶最大緩衝區容量。

NumberOfSubscribers

回傳目前訂閱人數。

PeerReference

取得這個 Java 節點的 JNI 物件參考。

(繼承來源 Object)
Subscribers

回傳現有訂閱者清單,用於監控與追蹤,而非用於呼叫 Flow.Subscriber 訂閱者的方法。

ThresholdClass

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

ThresholdType

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

方法

名稱 Description
Clone()

建立並傳回這個 對象的複本。

(繼承來源 Object)
Close()

除非已關閉,否則會向現有訂閱者發出 Flow.Subscriber#onComplete() onComplete 訊號,並禁止後續嘗試發佈。

CloseExceptionally(Throwable)

除非已關閉,否則會向現有訂閱者發出 Flow.Subscriber#onError(Throwable) onError 錯誤訊號,並禁止後續發佈嘗試。

Construct(JniObjectReference, JniObjectReferenceOptions)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
Consume(IConsumer)

使用給定的消費者函數處理所有已發佈的項目。

Dispose()

釋放該 Java 節點所持有的資源。

(繼承來源 Object)
Dispose(Boolean)

釋放該 Java 節點所持有的資源。

(繼承來源 Object)
DisposeUnlessReferenced()

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
Equals(Object)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
Equals(Object)

指出其他物件是否「等於」這個物件。

(繼承來源 Object)
EstimateMaximumLag()

回傳所有現有訂閱者中尚未生產但尚未消費的最大物品數量估計值。

EstimateMinimumDemand()

回傳所有現有訂閱者中已請求 Flow.Subscription#request(long) request但尚未生產的最低項目數量估計值。

GetHashCode()

傳回此物件的雜湊碼值。

(繼承來源 Object)
IsSubscribed(Flow+ISubscriber)

如果該訂閱者目前是訂閱者,則會回傳為真。

JavaFinalize()
已淘汰.

當垃圾收集決定不再參考物件時,垃圾收集行程在 物件上呼叫。

(繼承來源 Object)
Notify()

喚醒正在等候此物件監視器的單一線程。

(繼承來源 Object)
NotifyAll()

喚醒正在等候此物件監視器的所有線程。

(繼承來源 Object)
Offer(Object, IBiPredicate)

若可能,透過非同步調用該 Flow.Subscriber#onNext(Object) onNext 方法,將該項目發佈給每位現有訂閱者。

Offer(Object, Int64, TimeUnit, IBiPredicate)

如果可能,透過非同步呼叫該 Flow.Subscriber#onNext(Object) onNext 方法,將指定項目發佈給每個當前訂閱者,當任何訂閱的資源無法使用時,阻塞至指定的逾時或呼叫執行緒被中斷,屆時呼叫該處理程序(若非空)會被呼叫,若回傳為真,則重新嘗試一次。

SetHandle(IntPtr, JniHandleOwnership)

設定 Handle 屬性。

(繼承來源 Object)
SetPeerReference(JniObjectReference, JniObjectReferenceOptions)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
Submit(Object)

透過非同步呼叫該 Flow.Subscriber#onNext(Object) onNext 方法,在任何訂閱者資源無法使用時,將指定項目發佈給每位現有訂閱者,且不間斷地封鎖。

Subscribe(Flow+ISubscriber)

除非已經訂閱,否則會新增該訂閱者。

ToArray<T>()

從這個 Java 陣列包裝器建立一個受管理陣列。

(繼承來源 Object)
ToString()

傳回物件的字串表示。

(繼承來源 Object)
UnregisterFromRuntime()

將此 Java 節點從互通執行時中取消註冊。

(繼承來源 Object)
Wait()

讓目前線程等候直到喚醒為止,通常是藉由em <notified/em>或<em>interrupted</em> 來喚醒它。<>

(繼承來源 Object)
Wait(Int64, Int32)

讓目前的線程等到喚醒為止,通常是因為 <em>notified</em> 或 <em>interrupted</em>,或直到經過一定數量的實時為止。

(繼承來源 Object)
Wait(Int64)

讓目前的線程等到喚醒為止,通常是因為 <em>notified</em> 或 <em>interrupted</em>,或直到經過一定數量的實時為止。

(繼承來源 Object)

明確介面實作

名稱 Description
IJavaPeerable.Disposed()

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.Finalized()

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.JniObjectReferenceControlBlock

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.SetJniIdentityHashCode(Int32)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.SetPeerReference(JniObjectReference)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

(繼承來源 JavaObject)
IJavaPeerable.UnregisterFromRuntime()

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

擴充方法

名稱 Description
GetJniTypeName(IJavaPeerable)

取得實例 self類型的 JNI 名稱。

JavaAs<TResult>(IJavaPeerable)

試著強制self輸入 TResult,檢查 強制在 Java 端是否有效。

JavaCast<TResult>(IJavaObject)

執行 Android 執行時間檢查的類型轉換。

JavaCast<TResult>(IJavaObject)

一個 Flow.Publisher 非同步地將提交的(非空)項目發送給現有訂閱者,直到訂閱結束。

TryJavaCast<TResult>(IJavaPeerable, TResult)

試著強制self輸入 TResult,檢查 強制在 Java 端是否有效。

適用於