public class MultipleBroadcasterMessageHandler extends AbstractReactorMessageHandler
MessageHandler のアイテムを一度に配信するように適応します。メッセージが委譲される特定のストリームは、partitionExpression 値によって決定されます。プロセッサーで inputStream のスケジュールを変更しない限り、partitionExpression が異なるメッセージバスディスパッチャースレッドで配信されたメッセージを同じストリームにマップしないようにする必要があります。これは、基盤となる Broadcaster の使用によるものです。たとえば、式 T(java.lang.Thread).currentThread().getId() を使用すると、現在のディスパッチャースレッド ID が Stream のインスタンスにマップされます。Kafka パーティションごとに Stream が必要な場合は、MessageBus ディスパッチャースレッドが各パーティションで同じになるため、式 header['kafka_partition_id'] を使用できます。 partitionExpression 値にマップされたストリームにエラーがある場合、または完了した場合は、次に消費されるメッセージが同じ partitionExpression 値にマップされたときに再作成されます。すべてのエラー処理はプロセッサー実装の責任です。AbstractReactorMessageHandler.ChannelForwardingSubscriberlogger, processor| コンストラクターと説明 |
|---|
MultipleBroadcasterMessageHandler(Processor processor, java.lang.String partitionExpression) 処理を委譲するリアクターベースのプロセッサーとパーティション式を指定して、新しい MessageHandler を構築します。 |
| 修飾子と型 | メソッドと説明 |
|---|---|
void | destroy() |
protected void | handleMessageInternal(org.springframework.messaging.Message<?> message) |
protected void | onInit() |
void | setIntegrationEvaluationContext(org.springframework.expression.EvaluationContext evaluationContext) |
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 MultipleBroadcasterMessageHandler(Processor processor, java.lang.String partitionExpression)
processor - ストリームベースのリアクタープロセッサー public void setIntegrationEvaluationContext(org.springframework.expression.EvaluationContext evaluationContext)
protected void onInit()
throws java.lang.Exceptionorg.springframework.integration.handler.AbstractMessageProducingHandler の onInit java.lang.Exceptionprotected void handleMessageInternal(org.springframework.messaging.Message<?> message)
org.springframework.integration.handler.AbstractMessageHandler の handleMessageInternal public void destroy()
throws java.lang.Exceptionjava.lang.Exception