SubmissionPublisher 類別
定義
重要
部分資訊涉及發行前產品,在發行之前可能會有大幅修改。 Microsoft 對此處提供的資訊,不做任何明確或隱含的瑕疵擔保。
一個 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
- 繼承
- 屬性
- 實作
備註
一個 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 |
| SubmissionPublisher(IExecutor, Int32, IBiConsumer) |
使用給定的執行器建立新的 SubmissionPublisher,進行非同步傳遞給訂閱者,每個訂閱者的最大緩衝區大小為指定,若非空,則當任何訂閱者在方法 |
| SubmissionPublisher(IExecutor, Int32) |
使用給定的執行器建立新的 SubmissionPublisher,進行非同步傳遞給訂閱者,每個訂閱者擁有最大緩衝區大小,且方法中不包含訂閱者例外 |
| SubmissionPublisher(IntPtr, JniHandleOwnership) |
一個 |
屬性
| 名稱 | Description |
|---|---|
| Class |
傳回這個 |
| ClosedException |
回傳與 |
| Executor |
回傳用於非同步傳送的執行者。 |
| Handle |
基礎Android實例的句柄。 (繼承來源 Object) |
| HasSubscribers |
如果該出版商有任何訂閱者,則回傳為真。 |
| IsClosed |
如果該出版社不接受投稿,則會回傳為真。 |
| JniIdentityHashCode |
取得由互通執行時指派給此 Java 對等端的身份雜湊碼。 (繼承來源 Object) |
| JniManagedPeerState |
一個 |
| JniPeerMembers |
一個 |
| MaxBufferCapacity |
回傳每位用戶最大緩衝區容量。 |
| NumberOfSubscribers |
回傳目前訂閱人數。 |
| PeerReference |
取得這個 Java 節點的 JNI 物件參考。 (繼承來源 Object) |
| Subscribers |
回傳現有訂閱者清單,用於監控與追蹤,而非用於呼叫 |
| ThresholdClass |
一個 |
| ThresholdType |
一個 |
方法
| 名稱 | Description |
|---|---|
| Clone() |
建立並傳回這個 對象的複本。 (繼承來源 Object) |
| Close() |
除非已關閉,否則會向現有訂閱者發出 |
| CloseExceptionally(Throwable) |
除非已關閉,否則會向現有訂閱者發出 |
| Construct(JniObjectReference, JniObjectReferenceOptions) |
一個 |
| Consume(IConsumer) |
使用給定的消費者函數處理所有已發佈的項目。 |
| Dispose() |
釋放該 Java 節點所持有的資源。 (繼承來源 Object) |
| Dispose(Boolean) |
釋放該 Java 節點所持有的資源。 (繼承來源 Object) |
| DisposeUnlessReferenced() |
一個 |
| Equals(Object) |
一個 |
| Equals(Object) |
指出其他物件是否「等於」這個物件。 (繼承來源 Object) |
| EstimateMaximumLag() |
回傳所有現有訂閱者中尚未生產但尚未消費的最大物品數量估計值。 |
| EstimateMinimumDemand() |
回傳所有現有訂閱者中已請求 |
| GetHashCode() |
傳回此物件的雜湊碼值。 (繼承來源 Object) |
| IsSubscribed(Flow+ISubscriber) |
如果該訂閱者目前是訂閱者,則會回傳為真。 |
| JavaFinalize() |
已淘汰.
當垃圾收集決定不再參考物件時,垃圾收集行程在 物件上呼叫。 (繼承來源 Object) |
| Notify() |
喚醒正在等候此物件監視器的單一線程。 (繼承來源 Object) |
| NotifyAll() |
喚醒正在等候此物件監視器的所有線程。 (繼承來源 Object) |
| Offer(Object, IBiPredicate) |
若可能,透過非同步調用該 |
| Offer(Object, Int64, TimeUnit, IBiPredicate) |
如果可能,透過非同步呼叫該 |
| SetHandle(IntPtr, JniHandleOwnership) |
設定 Handle 屬性。 (繼承來源 Object) |
| SetPeerReference(JniObjectReference, JniObjectReferenceOptions) |
一個 |
| Submit(Object) |
透過非同步呼叫該 |
| 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() |
一個 |
| IJavaPeerable.Finalized() |
一個 |
| IJavaPeerable.JniObjectReferenceControlBlock |
一個 |
| IJavaPeerable.SetJniIdentityHashCode(Int32) |
一個 |
| IJavaPeerable.SetJniManagedPeerState(JniManagedPeerStates) |
一個 |
| IJavaPeerable.SetPeerReference(JniObjectReference) |
一個 |
| IJavaPeerable.UnregisterFromRuntime() |
一個 |
擴充方法
| 名稱 | Description |
|---|---|
| GetJniTypeName(IJavaPeerable) |
取得實例 |
| JavaAs<TResult>(IJavaPeerable) |
試著強制 |
| JavaCast<TResult>(IJavaObject) |
執行 Android 執行時間檢查的類型轉換。 |
| JavaCast<TResult>(IJavaObject) |
一個 |
| TryJavaCast<TResult>(IJavaPeerable, TResult) |
試著強制 |