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