public class BroadcasterMessageHandler extends AbstractReactorMessageHandler
Stream に処理を委譲することで、MessageHandler の配信時にアイテムを適応させます。プロセッサーの outputStream は、メッセージを作成し、それを出力チャネルに送信するために使用されます。入力チャネルと出力チャネルが MessageBus に接続されている場合、onNext の呼び出しを介して入力ストリームに配信されたデータは、メッセージバスのディスパッチャースレッドで呼び出され、出力チャネルにメッセージを送信すると、メッセージバス上の IO 操作が行われます。実装では、非同期ディスパッチを備えた RingBufferProcessor を使用します。これには、ストリームの状態を、onNext を呼び出すすべての受信ディスパッチャースレッド間で共有できるという利点があります。欠点は、処理と出力チャネルへの送信がディスパッチャースレッドの 1 つでシリアルに実行されることです。このハンドラーを使用すると、データを処理するときに非常に自然な最初のエクスペリエンスが得られます。たとえば、ストリーム http | react-processor | log があり、reactor-processor が buffer(5) を実行してから単一の値を生成するとします。http ソースに 10 個のメッセージを送信すると、使用されるディスパッチャースレッドの数に関係なく、ログには 2 つのメッセージが生成されます。プロセッサーから outputStream を返す前に、dispatchOn またはその他のスイッチ (http://projectreactor.io/docs/reference/#streams-multithreading)) を明示的に呼び出すことで、出力チャネルへの送信を行う outputStream サブスクライバが使用するスレッドを変更できます。複数のストリームにまたがるディスパッチャースレッドでの同時実行には、MultipleBroadcasterMessageHandler を使用します。すべてのエラー処理はプロセッサー実装の責任です。AbstractReactorMessageHandler.ChannelForwardingSubscriberlogger, processor| コンストラクターと説明 |
|---|
BroadcasterMessageHandler(Processor processor) 処理を委譲するリアクターベースのプロセッサーを指定して、新しい BroadcasterMessageHandler を構築します。 |
| 修飾子と型 | メソッドと説明 |
|---|---|
void | destroy() |
protected void | handleMessageInternal(org.springframework.messaging.Message<?> message) |
protected void | onInit() |
getEnvironment, getInputType, getRingBufferSize, getStopTimeout, invokeProcessor, setRingBufferSize, setStopTimeoutgetOutputChannel, produceOutput, sendOutputs, setOutputChannel, setOutputChannelName, setSendTimeout, shouldCopyRequestHeaders, shouldSplitOutputconfigureMetrics, getActiveCount, getActiveCountLong, getComponentType, getDuration, getErrorCount, getErrorCountLong, getHandleCount, getHandleCountLong, getManagedName, getManagedType, getMaxDuration, getMeanDuration, getMinDuration, getOrder, getStandardDeviationDuration, handleMessage, isCountsEnabled, isLoggingEnabled, isStatsEnabled, reset, setCountsEnabled, setLoggingEnabled, setManagedName, setManagedType, setOrder, setShouldTrack, setStatsEnabledafterPropertiesSet, extractTypeIfPossible, getApplicationContext, getApplicationContextId, getBeanFactory, getChannelResolver, getComponentName, getConversionService, getIntegrationProperties, getIntegrationProperty, getMessageBuilderFactory, getTaskScheduler, setApplicationContext, setBeanFactory, setBeanName, setChannelResolver, setComponentName, setConversionService, setMessageBuilderFactory, setTaskScheduler, toStringpublic BroadcasterMessageHandler(Processor processor)
processor - ストリームベースのリアクタープロセッサー protected void handleMessageInternal(org.springframework.messaging.Message<?> message)
throws java.lang.Exceptionorg.springframework.integration.handler.AbstractMessageHandler の handleMessageInternal java.lang.Exceptionpublic void destroy()
throws java.lang.Exceptionjava.lang.Exceptionprotected void onInit()
throws java.lang.Exceptionorg.springframework.integration.handler.AbstractMessageProducingHandler の onInit java.lang.Exception