このバージョンはまだ開発中であり、まだ安定しているとは見なされていません。最新の安定バージョンについては、Spring Integration 7.1.1 を使用してください! |
Redis サポート
Spring Integration 2.1 は、Redis (英語) のサポートを導入しました: 「オープンソースの高度な Key-Value ストア」。このサポートは、Redis ベースの MessageStore と、Redis が PUBLISH、SUBSCRIBE、UNSUBSCRIBE (英語) コマンドを介してサポートするパブリッシュ / サブスクライブメッセージングアダプターの形式で提供されます。
この依存関係はプロジェクトに必要です:
Maven
Gradle
<dependency>
<groupId>org.springframework.integration</groupId>
<artifactId>spring-integration-redis</artifactId>
<version>7.2.0-M1</version>
</dependency>
implementation "org.springframework.integration:spring-integration-redis:7.2.0-M1"
Redis クライアント依存関係を含める必要があります (例: Lettuce (英語) )。
Redis をダウンロード、インストール、実行するには、Redis ドキュメント (英語) を参照してください。
Redis への接続
Redis とのやり取りを開始するには、まず接続を確立する必要があります。Spring Integration は、別の Spring プロジェクトである Spring Data Redis [GitHub] (英語) のサポートを利用しています。Spring Data Redis [GitHub] (英語) は、Spring の典型的な構成要素である ConnectionFactory と Template を提供しています。これらの抽象化により、複数の Redis クライアント Java API との統合が簡素化されます。現在、Spring Data と Redis は Jedis [GitHub] (英語) と Lettuce (英語) をサポートしています。
RedisConnectionFactory を使用する
Spring Data/Redis の RedisConnectionFactory は、Redis との接続を管理するための高レベルの抽象化です。以下のリストはインターフェース定義を示しています。
public interface RedisConnectionFactory extends PersistenceExceptionTranslator {
/**
* Provides a suitable connection for interacting with Redis.
* @return connection for interacting with Redis.
*/
RedisConnection getConnection();
} 次の例は、Java で LettuceConnectionFactory を作成する方法を示しています。
LettuceConnectionFactory cf = new LettuceConnectionFactory();
cf.afterPropertiesSet(); 次の例は、Spring の XML 構成で LettuceConnectionFactory を作成する方法を示しています。
<bean id="redisConnectionFactory"
class="o.s.data.redis.connection.lettuce.LettuceConnectionFactory">
<property name="port" value="7379" />
</bean>RedisConnectionFactory の実装は、ポートやホストなどの一連のプロパティを提供します。RedisConnectionFactory のインスタンスが存在すると、RedisTemplate を作成できます。
RedisTemplate を使用する
Spring の他のテンプレートクラス(JdbcTemplate や JmsTemplate など)と同様に、RedisTemplate は Redis データアクセスコードを簡素化するヘルパークラスです。RedisTemplate およびそのバリエーション(StringRedisTemplate など)の詳細については、Spring Data Redis ドキュメントを参照してください。
次の例は、Java で RedisTemplate のインスタンスを作成する方法を示しています。
RedisTemplate rt = new RedisTemplate<String, Object>();
rt.setConnectionFactory(redisConnectionFactory); 次の例は、Spring の XML 構成で RedisTemplate のインスタンスを作成する方法を示しています。
<bean id="redisTemplate"
class="org.springframework.data.redis.core.RedisTemplate">
<property name="connectionFactory" ref="redisConnectionFactory"/>
</bean>Redis でメッセージング
で記述されていたように導入、Redis はその PUBLISH、SUBSCRIBE、UNSUBSCRIBE コマンドによってメッセージング、パブリッシュサブスクライブのサポートを提供します。JMS および AMQP と同様に、Spring Integration は、Redis を介してメッセージを送受信するためのメッセージチャネルとアダプターを提供します。
Redis パブリッシュ / サブスクライブチャネル
JMS と同様に、プロデューサーとコンシューマーの両方が同じアプリケーションの一部となり、同じプロセス内で実行される場合があります。これは、受信チャネルアダプターと送信チャネルアダプターのペアを使用することで実現できます。しかし、Spring Integration の JMS サポートと同様に、このユースケースにはよりシンプルな方法があります。代わりに、次の例に示すように、パブリッシュ / サブスクライブチャネルを使用できます。
<int-redis:publish-subscribe-channel id="redisChannel" topic-name="si.test.topic"/>publish-subscribe-channel は、メインの Spring Integration 名前空間の通常の <publish-subscribe-channel/> 要素とほぼ同様に動作します。任意のエンドポイントの input-channel 属性と output-channel 属性の両方から参照できます。違いは、このチャネルが Redis トピック名(topic-name 属性で指定される String 値)によってサポートされていることです。ただし、JMS とは異なり、このトピックは事前に作成する必要はなく、Redis によって自動作成される必要もありません。Redis では、トピックはアドレスのロールを果たす単純な String 値です。プロデューサーとコンシューマーは、同じ String 値をトピック名として使用して通信できます。このチャネルへの単純なサブスクリプションは、プロデューシングエンドポイントとコンシューマーエンドポイント間で非同期のパブリッシュ / サブスクライブメッセージングが可能になることを意味します。ただし、単純な Spring Integration <channel/> 要素内に <queue/> 要素を追加することで作成される非同期メッセージチャネルとは異なり、メッセージはメモリ内キューに格納されません。代わりに、これらのメッセージは Redis を介して渡され、永続性とクラスタリングのサポート、他の非 Java プラットフォームとの相互運用性に依存できます。
Redis 受信チャネルアダプター
Redis 受信チャネルアダプター(RedisInboundChannelAdapter)は、他の受信アダプターと同じ方法で、受信 Redis メッセージを Spring メッセージに適合させます。プラットフォーム固有のメッセージ(この場合は Redis)を受信し、MessageConverter 戦略を使用して Spring メッセージに変換します。次の例は、Redis 受信チャネルアダプターを構成する方法を示しています。
Java DSL
Java
XML
@Bean
RedisConnectionFactory connectionFactory() {
return RedisContainerTest.connectionFactory();
}
@Bean
MessageConverter testConverter() {
return new SimpleMessageConverter();
}
@Bean
IntegrationFlow inboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from(Redis
.inboundChannelAdapter(redisConnectionFactory)
.topics(TOPIC_FOR_INBOUND_CHANNEL_ADAPTER)
.messageConverter(testConverter()))
.channel(c -> c.queue("inboundChannelAdapterQueueChannel"))
.get();
}
@Bean
RedisConnectionFactory connectionFactory() {
return RedisContainerTest.connectionFactory();
}
@Bean
MessageConverter testConverter() {
return new SimpleMessageConverter();
}
@Bean
RedisInboundChannelAdapter inboundChannelAdapter(RedisConnectionFactory connectionFactory) {
var adapter = new RedisInboundChannelAdapter(connectionFactory);
adapter.setTopics(redisChannelName);
adapter.setOutputChannel(channel);
adapter.setMessageConverter(testConverter());
return adapter;
}
<int-redis:inbound-channel-adapter id="redisAdapter"
topics="thing1, thing2"
channel="receiveChannel"
error-channel="testErrorChannel"
message-converter="testConverter" />
<bean id="redisConnectionFactory"
class="o.s.data.redis.connection.lettuce.LettuceConnectionFactory">
<property name="port" value="7379" />
</bean>
<bean id="testConverter" class="things.something.SampleMessageConverter" />
上記の例は、Redis 受信チャネルアダプターのシンプルながらも包括的な構成を示しています。この構成は、特定の Bean を自動検出するという、おなじみの Spring パラダイムに依存していることに注意してください。この場合、redisConnectionFactory は暗黙的にアダプターに注入されます。あるいは、connection-factory 属性を介してカスタム RedisConnectionFactory を注入することもできます。
また、上記の構成では、アダプターにカスタム MessageConverter が挿入されることに注意してください。このアプローチは、MessageConverter インスタンスを使用して Redis メッセージと Spring Integration メッセージペイロードを変換する JMS に似ています。デフォルトは SimpleMessageConverter です。
受信アダプターは、複数のトピック名をサブスクライブできます。topics 属性のコンマ区切り値のセットです。
バージョン 3.0 以降、受信アダプターは、既存の topics 属性に加えて、topic-patterns 属性を持つようになりました。この属性には、Redis トピックパターンのコンマ区切りセットが含まれています。Redis のパブリッシュ / サブスクライブに関する詳細については、Redis Pub/Sub (英語) を参照してください。
受信アダプターは、RedisSerializer を使用して Redis メッセージの本文を逆直列化できます。<int-redis:inbound-channel-adapter> の serializer 属性を空の文字列に設定すると、RedisSerializer プロパティの null 値が得られます。この場合、Redis メッセージの未加工の byte[] 本体は、メッセージペイロードとして提供されます。
バージョン 5.0 以降、<int-redis:inbound-channel-adapter> の task-executor 属性を使用することで、Executor インスタンスを受信アダプターに注入できるようになりました。また、受信した Spring Integration メッセージには、発行されたメッセージのソース(トピックまたはパターン)を示す RedisHeaders.MESSAGE_SOURCE ヘッダーが追加されました。これは、下流のルーティングロジックで使用できます。
Redis 送信チャネルアダプター
Redis 送信チャネルアダプターは、他の送信アダプターと同じ方法で、発信 Spring Integration メッセージを Redis メッセージに適合させます。Spring Integration メッセージを受信し、MessageConverter 戦略を使用してプラットフォーム固有のメッセージ(この場合は Redis)に変換します。次の例は、Redis 送信チャネルアダプターを構成する方法を示しています。
Java DSL
Java
XML
@Bean
RedisConnectionFactory connectionFactory() {
return RedisContainerTest.connectionFactory();
}
@Bean
MessageConverter testConverter() {
return new SimpleMessageConverter();
}
@Bean
IntegrationFlow outboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory) {
return flow -> flow
.handle(Redis.outboundChannelAdapter(redisConnectionFactory)
.topicExpression(new LiteralExpression(TOPIC_FOR_OUTBOUND_CHANNEL_ADAPTER))
.messageConverter(testConverter()));
}
@Bean
RedisConnectionFactory connectionFactory() {
return RedisContainerTest.connectionFactory();
}
@Bean
MessageConverter testConverter() {
return new SimpleMessageConverter();
}
@Bean
RedisPublishingMessageHandler outboundChannelAdapter(RedisConnectionFactory connectionFactory) {
var handler = new RedisPublishingMessageHandler(connectionFactory);
handler.setTopicExpression(new LiteralExpression(topic));
handler.setMessageConverter(testConverter());
return handler;
}
<int-redis:outbound-channel-adapter id="outboundAdapter"
channel="sendChannel"
topic="thing1"
message-converter="testConverter"/>
<bean id="redisConnectionFactory"
class="o.s.data.redis.connection.lettuce.LettuceConnectionFactory">
<property name="port" value="7379"/>
</bean>
<bean id="testConverter" class="things.something.SampleMessageConverter" />
この構成は、Redis 受信チャネルアダプターに対応しています。アダプターには、RedisConnectionFactory が暗黙的に挿入されます。これは、Bean 名として redisConnectionFactory で定義されています。この例には、オプションの(およびカスタムの) MessageConverter (testConverter Bean)も含まれています。
Spring Integration 3.0 以降、<int-redis:outbound-channel-adapter> は topic 属性の代替として、実行時にメッセージの Redis トピックを決定するための topic-expression 属性を提供しています。これらの属性は相互に排他的です。
Redis キュー受信チャネルアダプター
Spring Integration 3.0 では、Redis リストからメッセージを「ポップ」するためのキュー受信チャネルアダプターが導入されました。デフォルトでは「右ポップ」を使用しますが、「左ポップ」を使用するように設定することもできます。このアダプターはメッセージ駆動型で、内部リスナースレッドを使用し、ポーラーは使用しません。
次のリストは、queue-inbound-channel-adapter で使用可能なすべての属性を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow queueInboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from(Redis
.queueInboundChannelAdapter(
queueName, (6)
redisConnectionFactory) (5)
.id() (1)
.outputChannel() (2)
.autoStartup() (3)
.phase() (4)
.errorChannel() (7)
.serializer() (8)
.receiveTimeout() (9)
.recoveryInterval() (10)
.expectMessage() (11)
.taskExecutor() (12)
.rightPop() (13)
)
.get();
}
@Bean
RedisQueueMessageDrivenEndpoint queueInboundChannelAdapter() {
var adapter = new RedisQueueMessageDrivenEndpoint(
queueName, (6)
redisConnectionFactory (5)
);
adapter.setBeanName(); (1)
adapter.setOutputChannel(); (2)
adapter.setAutoStartup(); (3)
adapter.setPhase(); (4)
adapter.setErrorChannel(); (7)
adapter.setSerializer(); (8)
adapter.setReceiveTimeout(); (9)
adapter.setRecoveryInterval(); (10)
adapter.setExpectMessage(); (11)
adapter.setTaskExecutor(); (12)
adapter.setRightPop(); (13)
return adapter;
}
<int-redis:queue-inbound-channel-adapter id="" (1)
channel="" (2)
auto-startup="" (3)
phase="" (4)
connection-factory="" (5)
queue="" (6)
error-channel="" (7)
serializer="" (8)
receive-timeout="" (9)
recovery-interval="" (10)
expect-message="" (11)
task-executor="" (12)
right-pop=""/> (13)
| 1 | コンポーネント Bean 名。channel 属性が指定されていない場合、DirectChannel が作成され、この id 属性を Bean 名としてアプリケーションコンテキストに登録されます。この場合、エンドポイント自体は Bean 名 id + .adapter で登録されます。(Bean 名が thing1 であった場合、エンドポイントは thing1.adapter として登録されます。) |
| 2 | このエンドポイントから Message インスタンスを送信する MessageChannel。 |
| 3 | アプリケーションコンテキストの起動後にこのエンドポイントを自動的に起動するかどうかを指定するための SmartLifecycle 属性。デフォルトは true です。 |
| 4 | このエンドポイントが開始されるフェーズを指定する SmartLifecycle 属性。デフォルトは 0 です。 |
| 5 | RedisConnectionFactory Bean への参照。デフォルトは redisConnectionFactory です。 |
| 6 | Redis メッセージを取得するためにキューベースの「ポップ」操作が実行される Redis リストの名前。 |
| 7 | エンドポイントのリスニングタスクから例外を受信したときに ErrorMessage インスタンスを送信する MessageChannel。デフォルトでは、基礎となる MessagePublishingErrorHandler はアプリケーションコンテキストのデフォルト errorChannel を使用します。 |
| 8 | RedisSerializer Bean リファレンス。これは、「シリアライザーなし」を意味する空の文字列にすることができます。この場合、受信 Redis メッセージからの生の byte[] が Message ペイロードとして channel に送信されます。デフォルトでは JdkSerializationRedisSerializer です。 |
| 9 | キューからの Redis メッセージを待機する「ポップ」操作のタイムアウト(ミリ秒)。デフォルトは 1 秒です。 |
| 10 | "pop" 操作の例外の後、リスナータスクを再起動するまでにリスナータスクがスリープする時間(ミリ秒)。 |
| 11 | このエンドポイントが Redis キューからのデータに Message インスタンス全体が含まれることを期待するかどうかを指定します。この属性が true に設定されている場合、メッセージは何らかの形式のデシリアライゼーション(デフォルトでは JDK シリアライゼーション)を必要とするため、serializer を空の文字列にすることはできません。デフォルトは false です。 |
| 12 | Spring TaskExecutor (または標準 JDK 1.5 + Executor)Bean への参照。これは、基礎となるリスニングタスクに使用されます。デフォルトは SimpleAsyncTaskExecutor です。 |
| 13 | このエンドポイントが Redis リストからメッセージを読み取るために「右ポップ」(true の場合)または「左ポップ」(false の場合)を使用するかどうかを指定します。true の場合、Redis リストは、デフォルトの Redis キュー送信チャネルアダプターと共に使用されると、FIFO キューとして機能します。false に設定すると、「右プッシュ」でリストに書き込むソフトウェアで使用したり、スタックのようなメッセージの順序を実現したりできます。デフォルトは true です。バージョン 4.3 以降。 |
task-executor は複数のスレッドで処理するように構成する必要があります。そうしないと、エラー発生後に RedisQueueMessageDrivenEndpoint がリスナータスクの再起動を試みたときにデッドロックが発生する可能性があります。errorChannel を使用してこれらのエラーを処理し、再起動を回避することは可能ですが、アプリケーションがデッドロック状態に陥らないようにすることが推奨されます。TaskExecutor の実装については、Spring Framework 参考マニュアルを参照してください。 |
Redis キュー送信チャネルアダプター
Spring Integration 3.0 では、Spring Integration メッセージを Redis リストに「プッシュ」するためのキュー送信チャネルアダプターが導入されました。デフォルトでは「左プッシュ」が使用されますが、「右プッシュ」に変更することもできます。以下のリストは、Redis queue-outbound-channel-adapter で利用可能なすべての属性を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow queueOutboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from("inputChannel") (2)
.handle(Redis.queueOutboundChannelAdapter(
queueName (4)
or queueExpression, (5)
redisConnectionFactory) (3)
.serializer() (6)
.extractPayload() (7)
.leftPush() (8)
, e->e.id()) (1)
.get();
}
@Bean
@ServiceActivator(inputChannel = "channel") (2)
RedisQueueOutboundChannelAdapter queueOutboundChannelAdapter() {
var adapter = new RedisQueueOutboundChannelAdapter(
queueName (4)
or queueExpression, (5)
redisConnectionFactory (3)
);
adapter.setBeanName(); (1)
adapter.setSerializer(); (6)
adapter.setExtractPayload(); (7)
adapter.setLeftPush(); (8)
return adapter;
}
<int-redis:queue-outbound-channel-adapter id="" (1)
channel="" (2)
connection-factory="" (3)
queue="" (4)
queue-expression="" (5)
serializer="" (6)
extract-payload="" (7)
left-push=""/> (8)
| 1 | コンポーネント Bean 名。channel 属性が指定されていない場合、DirectChannel が作成され、この id 属性を Bean 名としてアプリケーションコンテキストに登録されます。この場合、エンドポイントは id に .adapter を加えた Bean 名で登録されます。(Bean 名が thing1 であった場合、エンドポイントは thing1.adapter として登録されます。) |
| 2 | このエンドポイントが Message インスタンスを受信する MessageChannel。 |
| 3 | RedisConnectionFactory Bean への参照。デフォルトは redisConnectionFactory です。 |
| 4 | Redis メッセージを送信するためにキューベースの「プッシュ」操作が実行される Redis リストの名前。この属性は queue-expression と相互に排他的です。 |
| 5 | Redis リストの名前を決定する SpEL Expression。実行時に受信 Message を #root 変数として使用します。この属性は queue と相互に排他的です。 |
| 6 | RedisSerializer Bean リファレンス。デフォルトは JdkSerializationRedisSerializer です。ただし、String ペイロードの場合、serializer 参照が提供されていない場合は、StringRedisSerializer が使用されます。 |
| 7 | このエンドポイントがペイロードのみを送信するか、Message 全体を Redis キューに送信するかを指定します。デフォルトは true です。 |
| 8 | このエンドポイントが Redis リストにメッセージを書き込むために「左プッシュ」(true の場合)または「右プッシュ」(false の場合)を使用するかどうかを指定します。true の場合、Redis リストは、デフォルトの Redis キュー受信チャネルアダプターと共に使用されると、FIFO キューとして機能します。false に設定すると、「左ポップ」でリストから読み取るソフトウェアで使用したり、スタックのようなメッセージの順序を実現したりできます。デフォルトは true です。バージョン 4.3 以降。 |
Redis アプリケーションイベント
Spring Integration 3.0 以降、Redis モジュールは IntegrationEvent の実装を提供しており、これは org.springframework.context.ApplicationEvent です。RedisExceptionEvent は、Redis 操作からの例外をカプセル化します(エンドポイントがイベントの「ソース」となります)。たとえば、<int-redis:queue-inbound-channel-adapter/> は BoundListOperations.rightPop 操作からの例外をキャッチした後、これらのイベントを発行します。例外は、汎用の org.springframework.data.redis.RedisSystemException または org.springframework.data.redis.RedisConnectionFailureException のいずれかです。これらのイベントを <int-event:inbound-channel-adapter/> で処理することは、バックグラウンドの Redis タスクの問題を特定し、管理アクションを実行できます。
Redis メッセージストア
エンタープライズ統合パターン(EIP)の書籍に従って、メッセージストア (英語) はメッセージを永続化します。これは、信頼性が懸念される状況において、メッセージをバッファリングする機能を持つコンポーネント(アグリゲータ、リシーケンサなど)を扱う際に役立ちます。Spring Integration では、MessageStore 戦略は、EIP でも説明されているクレームチェック (英語) パターンの基盤も提供します。
Spring Integration の Redis モジュールは RedisMessageStore を提供します。次の例は、アグリゲーターでそれを使用する方法を示しています。
<bean id="redisMessageStore" class="o.s.i.redis.store.RedisMessageStore">
<constructor-arg ref="redisConnectionFactory"/>
</bean>
<int:aggregator input-channel="inputChannel" output-channel="outputChannel"
message-store="redisMessageStore"/> 上記の例は Bean 構成であり、コンストラクター引数として RedisConnectionFactory を想定しています。
デフォルトでは、RedisMessageStore は Java シリアライゼーションを使用してメッセージをシリアライズします。ただし、異なるシリアライゼーション手法(JSON など)が必要な場合は、RedisMessageStore の valueSerializer プロパティにカスタムシリアライザーを設定できます。
フレームワークは、Message インスタンスと MessageHeaders インスタンス用の Jackson シリアライザーとデシリアライザーの実装 (それぞれ MessageJsonDeserializer と MessageHeadersJsonSerializer) を提供します。これらは、ObjectMapper の SimpleModule オプションで構成する必要があります。さらに、直列化された複雑なオブジェクトごとに型情報を追加するには、ObjectMapper で enableDefaultTyping を設定する必要があります。その type 情報は、デシリアライズ中に使用されます。フレームワークは、JacksonMessagingUtils.messagingAwareMapper() というユーティリティメソッドを提供します。これには、前述のすべてのプロパティとシリアライザーがすでに用意されています。このユーティリティメソッドには、セキュリティの脆弱性を回避するためにデシリアライズする Java パッケージを制限する trustedPackages 引数が付属しています。デフォルトの信頼できるパッケージ: java.util、java.lang、org.springframework.messaging.support、org.springframework.integration.support、org.springframework.integration.message、org.springframework.integration.store。RedisMessageStore で JSON 直列化を管理するには、次のような構成を適用する必要があります。
RedisMessageStore store = new RedisMessageStore(redisConnectionFactory);
ObjectMapper mapper = JacksonMessagingUtils.messagingAwareMapper();
RedisSerializer<Object> serializer = new GenericJackson3JsonRedisSerializer(mapper);
store.setValueSerializer(serializer); バージョン 4.3.12 から、RedisMessageStore は prefix オプションをサポートし、同じ Redis サーバー上のストアのインスタンスを区別できるようにします。
Redis チャネルメッセージストア
前述の RedisMessageStore は、各グループを単一のキー(グループ ID)の値として保持します。QueueChannel は永続化に使用できますが、その目的のために専用の RedisChannelMessageStore が提供されています(バージョン 4.0 以降)。このストアは、各チャネルに LIST、メッセージ送信時に LPUSH、メッセージ受信時に RPOP を使用します。デフォルトでは、このストアは JDK シリアライゼーションも使用しますが、前述のように、値シリアライザ用に変更できます。
一般的な RedisMessageStore ではなく、ストアバッキングチャネルを使用することをお勧めします。次の例では、Redis メッセージストアを定義し、キューを持つチャネルで使用しています。
<bean id="redisMessageStore" class="o.s.i.redis.store.RedisChannelMessageStore">
<constructor-arg ref="redisConnectionFactory"/>
</bean>
<int:channel id="somePersistentQueueChannel">
<int:queue message-store="redisMessageStore"/>
<int:channel> データの保存に使用されるキーの形式は、<storeBeanName>:<channelId> (前の例では redisMessageStore:somePersistentQueueChannel)です。
さらに、サブクラス RedisChannelPriorityMessageStore も提供されています。これを QueueChannel と併用すると、メッセージは(FIFO)優先度順に受信されます。標準の IntegrationMessageHeaderAccessor.PRIORITY ヘッダーを使用し、優先度値(0 - 9)をサポートします。その他の優先度のメッセージ(および優先度のないメッセージ)は、優先度のあるメッセージの後に FIFO 順に取得されます。
これらのストアは BasicMessageGroupStore のみを実装し、MessageGroupStore は実装しません。それらは、QueueChannel のバックアップなどの状況にのみ使用できます。 |
Redis メタデータストア
Spring Integration 3.0 は、Redis ベースの新しい MetadataStore (Javadoc) (メタデータストアを参照)実装を導入しました。RedisMetadataStore は、アプリケーションの再起動後も MetadataStore の状態を維持するために使用できます。このような MetadataStore 実装は、次のようなアダプターで使用できます。
これらのアダプターに新しい RedisMetadataStore を使用するよう指示するには、metadataStore という名前の Spring Bean を宣言します。フィード受信チャネルアダプターとフィード受信チャネルアダプターの両方が、宣言された RedisMetadataStore を自動的に取得して使用します。次の例は、そのような Bean を宣言する方法を示しています。
<bean name="metadataStore" class="o.s.i.redis.store.metadata.RedisMetadataStore">
<constructor-arg name="connectionFactory" ref="redisConnectionFactory"/>
</bean>RedisMetadataStore は RedisProperties (Javadoc) によってサポートされています。それとの相互作用は BoundHashOperations (Javadoc) を使用します。これは、Properties ストア全体に対して key を必要とします。MetadataStore の場合、この key はリージョンのロールを果たします。これは、いくつかのアプリケーションが同じ Redis サーバーを使用する分散環境で役立ちます。デフォルトでは、この key の値は MetaData です。
バージョン 4.0 以降、このストアは ConcurrentMetadataStore を実装し、キーの値を保存または変更できるインスタンスが 1 つだけである複数のアプリケーションインスタンス間で確実に共有できるようにします。
アトミック性のための WATCH コマンドは現在サポートされていないため、RedisMetadataStore.replace() は Redis クラスターでは使用できません (たとえば、AbstractPersistentAcceptOnceFileListFilter 内)。 |
Redis ストア受信チャネルアダプター
Redis ストア受信チャネルアダプターは、Redis コレクションからデータを読み取り、Message ペイロードとして送信するポーリングコンシューマーです。次の例は、Redis ストア受信チャネルアダプターを構成する方法を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow storeInboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from(Redis
.storeInboundChannelAdapter(redisConnectionFactory, STORE_FOR_INBOUND_CHANNEL_ADAPTER)
.collectionType(CollectionType.LIST),
endpointConfigure -> endpointConfigure
.poller(Pollers.fixedDelay(2000))
.id("storeSourcePollingChannelAdapter"))
.channel(c -> c.queue("redisChannel"))
.get();
}
@Bean
@InboundChannelAdapter(channel = "redisChannel", poller = @Poller(fixedDelay = "2000"))
public RedisStoreMessageSource storeInboundChannelAdapter() {
var messageSource = new RedisStoreMessageSource(redisConnectionFactory, new LiteralExpression(STORE_FOR_INBOUND_CHANNEL_ADAPTER));
messageSource.setCollectionType(CollectionType.LIST);
return messageSource;
}
<int-redis:store-inbound-channel-adapter id="listAdapter"
connection-factory="redisConnectionFactory"
key="myCollection"
channel="redisChannel"
collection-type="LIST" >
<int:poller fixed-rate="2000" max-messages-per-poll="10"/>
</int-redis:store-inbound-channel-adapter>
上記の例は、store-inbound-channel-adapter 要素を使用して Redis ストア受信チャネルアダプターを構成し、次のようなさまざまな属性の値を提供する方法を示しています。
keyまたはkey-expression: 使用されているコレクションのキーの名前。collection-type: このアダプターでサポートされているコレクション型の列挙。サポートされているコレクションはLIST、SET、ZSET、PROPERTIES、MAPです。connection-factory:o.s.data.redis.connection.RedisConnectionFactoryのインスタンスへの参照。redis-template:o.s.data.redis.core.RedisTemplateのインスタンスへの参照。すべての受信 アダプターに共通するその他の属性 (「チャネル」など)。
redis-template と connection-factory は相互に排他的です。 |
デフォルトでは、アダプターは
|
key にリテラル値があるため、上記の例は比較的単純で静的です。場合によっては、何らかの条件に基づいて実行時にキーの値を変更する必要があります。そのような場合は、代わりに key-expression を使用します。ここで、指定する式は有効な SpEL 式であれば何でも構いません。
また、Redis コレクションから読み取られた、処理に成功したデータに対して、何らかの後処理を実行することもできます。たとえば、処理後に値を移動または削除するなどです。このようなロジックには、トランザクション同期機能を使用できます。次の例では、key-expression とトランザクション同期を使用しています。
Java DSL
Java
XML
@Bean
IntegrationFlow storeInboundChannelAdapterFlow(RedisConnectionFactory redisConnectionFactory,
TransactionSynchronizationFactory syncFactory) {
return IntegrationFlow.from(Redis
.storeInboundChannelAdapter(redisConnectionFactory, STORE_FOR_INBOUND_CHANNEL_ADAPTER)
.collectionType(CollectionType.ZSET),
endpointConfigure -> endpointConfigure
.poller(Pollers
.fixedDelay(1000)
.transactional(new PseudoTransactionManager())
.transactionSynchronizationFactory(syncFactory)
)
.autoStartup(false)
.id("storeSourcePollingChannelAdapter"))
.transform((GenericTransformer<RedisZSet<?>, Collection<?>>) source -> source.rangeByScore(1, 3))
.channel(c -> c.queue("storeInboundChannelAdapterOutputChannel"))
.get();
}
@Bean
QueueChannel afterCommitChannel() {
return new QueueChannel();
}
@Bean
TransactionSynchronizationProcessor syncProcessor(QueueChannel afterCommitChannel) {
var processor = new ExpressionEvaluatingTransactionSynchronizationProcessor();
SpelExpressionParser parser = new SpelExpressionParser();
processor.setAfterCommitExpression(parser.parseExpression("payload.removeByScore(1, 3)"));
processor.setAfterCommitChannel(afterCommitChannel);
return processor;
}
@Bean
TransactionSynchronizationFactory syncFactory(TransactionSynchronizationProcessor syncProcessor) {
return new DefaultTransactionSynchronizationFactory(syncProcessor);
}
@Bean
@InboundChannelAdapter(channel = "outputChannel", autoStartup = "false", poller = @Poller("myPoller"))
public RedisStoreMessageSource storeInboundChannelAdapter() {
var messageSource = new RedisStoreMessageSource(redisConnectionFactory, new LiteralExpression(STORE_FOR_INBOUND_CHANNEL_ADAPTER));
messageSource.setCollectionType(CollectionType.ZSET);
return messageSource;
}
@Bean
QueueChannel afterCommitChannel() {
return new QueueChannel();
}
@Bean
public PollerMetadata myPoller(TransactionSynchronizationFactory syncFactory) {
PollerMetadata poller = new PollerMetadata();
PeriodicTrigger periodicTrigger = new PeriodicTrigger(Duration.ofMillis(1000));
periodicTrigger.setFixedRate(false);
poller.setTrigger(periodicTrigger);
var txInterceptor = new TransactionInterceptorBuilder()
.transactionManager(new PseudoTransactionManager())
.build();
poller.setTransactionSynchronizationFactory(syncFactory);
poller.setAdviceChain(Collections.singletonList(txInterceptor));
return poller;
}
@Bean
TransactionSynchronizationProcessor syncProcessor(QueueChannel afterCommitChannel) {
var processor = new ExpressionEvaluatingTransactionSynchronizationProcessor();
SpelExpressionParser parser = new SpelExpressionParser();
processor.setAfterCommitExpression(parser.parseExpression("payload.removeByScore(1, 3)"));
processor.setAfterCommitChannel(afterCommitChannel);
return processor;
}
@Bean
TransactionSynchronizationFactory syncFactory(TransactionSynchronizationProcessor syncProcessor) {
return new DefaultTransactionSynchronizationFactory(syncProcessor);
}
<int-redis:store-inbound-channel-adapter id="zsetAdapterWithSingleScoreAndSynchronization"
connection-factory="redisConnectionFactory"
key-expression="'presidents'"
channel="otherRedisChannel"
auto-startup="false"
collection-type="ZSET">
<int:poller fixed-rate="1000" max-messages-per-poll="2">
<int:transactional synchronization-factory="syncFactory"/>
</int:poller>
</int-redis:store-inbound-channel-adapter>
<int:transaction-synchronization-factory id="syncFactory">
<int:after-commit expression="payload.removeByScore(18, 18)"/>
</int:transaction-synchronization-factory>
<bean id="transactionManager" class="o.s.i.transaction.PseudoTransactionManager"/>
transactional 要素を使用することで、ポーラーをトランザクション対応にすることができます。この要素は、たとえばフローの他の部分で JDBC が呼び出される場合など、実際のトランザクションマネージャーを参照できます。「実際の」トランザクションが存在しない場合は、代わりに o.s.i.transaction.PseudoTransactionManager を使用できます。o.s.i.transaction.PseudoTransactionManager は Spring の PlatformTransactionManager の実装であり、実際のトランザクションがない場合でも Redis アダプターのトランザクション同期機能を使用できます。
| これにより、Redis アクティビティ自体がトランザクションになりません。成功(コミット)または失敗(ロールバック)の前後にアクションの同期を取ることができます。 |
ポーラーがトランザクション対応になると、o.s.i.transaction.TransactionSynchronizationFactory のインスタンスを transactional 要素に追加できます。TransactionSynchronizationFactory は TransactionSynchronization のインスタンスを作成します。便宜上、デフォルトの SpEL ベースの TransactionSynchronizationFactory が公開されており、これを使用して SpEL 式を構成し、その実行をトランザクションと調整 (同期) できます。before-commit、after-commit、after-rollback の式が、評価結果 (ある場合) が送信されるチャネル (イベントの種類ごとに 1 つ) とともにサポートされています。子要素ごとに、expression 属性と channel 属性を指定できます。channel 属性のみが存在する場合、受信メッセージは特定の同期シナリオの一部としてそこに送信されます。expression 属性のみが存在し、式の結果が null 以外の値である場合は、結果をペイロードとして含むメッセージが生成され、デフォルトチャネル (NullChannel) に送信されて、ログ (DEBUG レベル) に表示されます。式の結果が null または void の場合、メッセージは生成されません。
RedisStoreMessageSource は、TransactionSynchronizationProcessor 実装からアクセスできるトランザクション IntegrationResourceHolder にバインドされた RedisStore インスタンスを持つ store 属性を追加します。
トランザクション同期の詳細については、トランザクションの同期を参照してください。
RedisStore 送信チャネルアダプター
RedisStore 送信チャネルアダプターを使用すると、次の例に示すように、メッセージペイロードを Redis コレクションに書き込むことができます。
Java DSL
Java
XML
@Bean
public IntegrationFlow storeOutboundChannelAdapterFlow(RedisConnectionFactory connectionFactory) {
return IntegrationFlow.from("requestChannel")
.handle(Redis.storeOutboundChannelAdapter(connectionFactory)
.collectionType(CollectionType.LIST)
.key("myCollection"))
.get();
}
@Bean
@ServiceActivator(inputChannel = "requestChannel")
public MessageHandler redisListAdapter(RedisConnectionFactory connectionFactory) {
var handler = new RedisStoreWritingMessageHandler(connectionFactory);
handler.setCollectionType(CollectionType.LIST);
handler.setKey("myCollection");
return handler;
}
<int-redis:store-outbound-channel-adapter id="redisListAdapter"
collection-type="LIST"
channel="requestChannel"
key="myCollection" />
上記の構成では、store-inbound-channel-adapter 要素を使用して Redis ストア送信チャネルアダプターを保存します。次のようなさまざまな属性の値を提供します。
keyまたはkey-expression: 使用されているコレクションのキーの名前。extract-payload-elements:true(デフォルト) に設定され、ペイロードが「複数値」オブジェクト (つまり、CollectionまたはMap) のインスタンスである場合は、"addAll" および "putAll" セマンティクスを使用して格納されます。それ以外の場合、falseに設定すると、ペイロードはその型に関係なく単一のエントリとして格納されます。ペイロードが「複数値」オブジェクトのインスタンスでない場合、この属性の値は無視され、ペイロードは常に単一のエントリとして格納されます。collection-type: このアダプターでサポートされているCollection型の列挙。サポートされているコレクションはLIST、SET、ZSET、PROPERTIES、MAPです。map-key-expression: 格納されているエントリのキーの名前を返す SpEL 式。collection-typeがMAPまたはPROPERTIESであり、"extract-payload-elements" が false の場合にのみ適用されます。connection-factory:o.s.data.redis.connection.RedisConnectionFactoryのインスタンスへの参照。redis-template:o.s.data.redis.core.RedisTemplateのインスタンスへの参照。すべての受信 アダプターに共通するその他の属性 (「チャネル」など)。
redis-template と connection-factory は相互に排他的です。 |
デフォルトでは、アダプターは StringRedisTemplate を使用します。これは、キー、値、ハッシュキー、ハッシュ値に StringRedisSerializer インスタンスを使用します。ただし、extract-payload-elements が false に設定されている場合、キーとハッシュキーに StringRedisSerializer インスタンス、値とハッシュ値に JdkSerializationRedisSerializer インスタンスを持つ RedisTemplate が使用されます。JDK シリアライザを使用する場合、値が実際にコレクションであるかどうかに関係なく、すべての値に Java シリアライゼーションが使用されることを理解することが重要です。値のシリアライゼーションをより細かく制御したい場合は、これらのデフォルトに依存せずに、カスタム RedisTemplate を提供することができます。 |
上記の例は、key 属性やその他の属性にリテラル値があるため、比較的単純で静的です。場合によっては、実行時に何らかの条件に基づいて値が動的に変更されることがあります。そのために、-expression に相当するもの(key-expression、map-key-expression など)が提供されています。式は、有効な任意の SpEL 式です。
Redis 送信コマンドゲートウェイ
Spring Integration 4.0 では、汎用的な RedisConnection#execute メソッドを使用して標準的な Redis コマンドを実行できるように、Redis コマンドゲートウェイが導入されました。以下のリストは、Redis 送信ゲートウェイで利用可能な属性を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow outboundGatewayFlow() {
return IntegrationFlow.from("request-channel") (1)
.handle(Redis
.outboundGateway(redisConnectionFactory (5)
or redisTemplate) (6)
.argumentsSerializer() (7)
.commandExpression("") (8)
.argumentsStrategy(new ExpressionArgumentsStrategy(
<argument-expressions>, (9)
<use-command-variable>) (10)
)
.argumentsStrategy() (11)
, e -> e
.requiresReply() (3)
.sendTimeout() (4)
)
.channel(c -> c.queue("replyChannel")) (2)
.get();
}
@Bean
@ServiceActivator(
inputChannel = "", (1)
outputChannel = "", (2)
requiresReply = "", (3)
sendTimeout = "") (4)
MessageHandler outboundGateway() {
RedisOutboundGateway gateway = new RedisOutboundGateway(
redisConnectionFactory (5)
or redisTemplate ); (6)
gateway.setArgumentsSerializer(); (7)
gateway.setCommandExpression(); (8)
var expressionArgumentStrategy = new ExpressionArgumentsStrategy(
<argument-expressions>, (9)
<use-command-variable>); (10)
gateway.setArgumentsStrategy(expressionArgumentStrategy);
gateway.setArgumentsStrategy(); (11)
return gateway;
}<int-redis:outbound-gateway
request-channel="" (1)
reply-channel="" (2)
requires-reply="" (3)
reply-timeout="" (4)
connection-factory="" (5)
redis-template="" (6)
arguments-serializer="" (7)
command-expression="" (8)
argument-expressions="" (9)
use-command-variable="" (10)
arguments-strategy="" /> (11)
| 1 | このエンドポイントが Message インスタンスを受信する MessageChannel。 |
| 2 | このエンドポイントが応答 Message インスタンスを送信する MessageChannel。 |
| 3 | この送信ゲートウェイが null 以外の値を返す必要があるかどうかを指定します。デフォルトは true です。Redis が null 値を返すと、ReplyRequiredException がスローされます。 |
| 4 | 応答メッセージが送信されるまで待機するタイムアウト(ミリ秒単位)。通常、キューベースの制限された応答チャネルに適用されます。 |
| 5 | RedisConnectionFactory Bean への参照。デフォルトは redisConnectionFactory です。redis-template 属性とは排他的です。 |
| 6 | RedisTemplate Bean への参照。connection-factory 属性とは排他的です。 |
| 7 | org.springframework.data.redis.serializer.RedisSerializer インスタンスへの参照。必要に応じて、各コマンド引数を byte[] に直列化するために使用されます。 |
| 8 | コマンドキーを返す SpEL 式。デフォルトは redis_command メッセージヘッダーです。null に評価してはなりません。 |
| 9 | コマンド引数として評価される、カンマ区切りの SpEL 式です。arguments-strategy 属性とは排他的です。どちらの属性も指定されていない場合は、payload がコマンド引数として使用されます。引数式は、可変数の引数をサポートするために "null" と評価されることがあります。 |
| 10 | argument-expressions が構成されている場合に、評価された Redis コマンド文字列を o.s.i.redis.outbound.ExpressionArgumentsStrategy の式評価コンテキストで #cmd 変数として使用できるようにするかどうかを指定する boolean フラグ。それ以外の場合、この属性は無視されます。 |
| 11 | o.s.i.redis.outbound.ArgumentsStrategy のインスタンスへの参照。argument-expressions 属性とは相互に排他的です。どちらの属性も指定されていない場合は、コマンド引数として payload が使用されます。 |
<int-redis:outbound-gateway> は、任意の Redis 演算を実行するための共通コンポーネントとして使用できます。次の例は、Redis 原子番号から増分値を取得する方法を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow outboundGatewayFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from("requestChannel")
.handle(Redis.outboundGateway(redisConnectionFactory)
.command("INCR"))
.channel(c -> c.queue("replyChannel"))
.get();
}@Bean
@ServiceActivator(inputChannel = "requestChannel", outputChannel = "replyChannel")
MessageHandler outboundGateway(RedisConnectionFactory redisConnectionFactory) {
RedisOutboundGateway gateway = new RedisOutboundGateway(redisConnectionFactory);
gateway.setCommandExpressionString("'INCR'");
return gateway;
}
<int-redis:outbound-gateway request-channel="requestChannel"
reply-channel="replyChannel"
command-expression="'INCR'"/>
Message ペイロードの名前は redisCounter である必要があり、これは org.springframework.data.redis.support.atomic.RedisAtomicInteger Bean 定義によって提供される場合があります。
RedisConnection#execute メソッドの戻り値は汎用の Object 型です。実際の結果はコマンドの種類によって異なります。例: MGET は List<byte[]> を返します。コマンド、その引数、結果の型の詳細については、Redis 仕様 (英語) を参照してください。
Redis キュー送信ゲートウェイ
Spring Integration では、リクエストと応答のシナリオを実行するために Redis キュー送信ゲートウェイが導入されました。会話 UUID を指定された queue にプッシュし、その UUID をキーとして持つ値を Redis リストにプッシュし、UUID と .reply のキーを持つ Redis リストからの応答を待ちます。インタラクションごとに異なる UUID が使用されます。次のリストは、Redis 送信ゲートウェイで使用可能な属性を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow queueOutboundGatewayFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from("requestChannel") (1)
.handle(Redis.queueOutboundGateway(
QUEUE_NAME, (6)
redisConnectionFactory) (5)
.serializer() (8)
.extractPayload() (9)
, e -> e
.requiresReply() (3)
.sendTimeout() (4)
.order()) (7)
.channel(c -> c.queue("replyChannel")) (2)
.get();
}
@Bean
@ServiceActivator(
inputChannel = "", (1)
outputChannel = "", (2)
requiresReply = "", (3)
sendTimeout = "" (4)
)
MessageHandler queueOutboundGateway(RedisConnectionFactory redisConnectionFactory) {
var gateway = new RedisQueueOutboundGateway(
QUEUE_NAME, (6)
redisConnectionFactory); (5)
gateway.setOrder(); (7)
gateway.setSerializer(); (8)
gateway.setExtractPayload(); (9)
return gateway;
}
<int-redis:queue-outbound-gateway
request-channel="" (1)
reply-channel="" (2)
requires-reply="" (3)
reply-timeout="" (4)
connection-factory="" (5)
queue="" (6)
order="" (7)
serializer="" (8)
extract-payload=""/> (9)
| 1 | このエンドポイントが Message インスタンスを受信する MessageChannel。 |
| 2 | このエンドポイントが応答 Message インスタンスを送信する MessageChannel。 |
| 3 | この送信ゲートウェイが非 null 値を返す必要があるかどうかを指定します。この値は、デフォルトでは false です。そうでない場合、Redis が null 値を返すと、ReplyRequiredException がスローされます。 |
| 4 | 応答メッセージが送信されるまで待機するタイムアウト(ミリ秒単位)。通常、キューベースの制限された応答チャネルに適用されます。 |
| 5 | RedisConnectionFactory Bean への参照。デフォルトは redisConnectionFactory です。'redis-template' 属性と相互に排他的です。 |
| 6 | 送信ゲートウェイが会話 UUID を送信する Redis リストの名前。 |
| 7 | 複数のゲートウェイが登録されている場合のこの送信ゲートウェイの順序。 |
| 8 | RedisSerializer Bean リファレンス。空の文字列にすることもできます。これは、「シリアライザなし」を意味します。この場合、受信 Redis メッセージからの生の byte[] は、Message ペイロードとして channel に送信されます。デフォルトでは、JdkSerializationRedisSerializer です。 |
| 9 | このエンドポイントがペイロードのみを処理するか、Redis キューに対するメッセージ全体を処理するかを指定します。デフォルトは true です。 |
Redis キュー受信ゲートウェイ
Spring Integration 4.1 は、リクエストと応答のシナリオを実行するために Redis キューの受信ゲートウェイを導入しました。提供された queue から会話 UUID をポップし、その UUID をキーとする値を Redis リストからポップし、UUID と .reply のキーを使用して Redis リストへの応答をプッシュします。次のリストは、Redis キューの受信ゲートウェイで使用可能な属性を示しています。
Java DSL
Java
XML
@Bean
IntegrationFlow queueInboundGatewayFlow(RedisConnectionFactory redisConnectionFactory) {
return IntegrationFlow.from("requestChannel") (1)
.handle(Redis.queueInboundGateway(
QUEUE_NAME, (6)
redisConnectionFactory) (5)
.taskExecutor() (3)
.replyTimeout() (4)
.phase() (7)
.serializer() (8)
.receiveTimeout() (9)
.extractPayload() (10)
.recoveryInterval()) (11)
.channel(c -> c.queue("replyChannel")) (2)
.get();
}
@Bean
public RedisQueueInboundGateway queueInboundGateway(RedisConnectionFactory connectionFactory) {
var gateway = new RedisQueueInboundGateway(
QUEUE_NAME, (6)
connectionFactory); (5)
gateway.setRequestChannel(); (1)
gateway.setReplyChannel(); (2)
gateway.setTaskExecutor(); (3)
gateway.setReplyTimeout(); (4)
gateway.setPhase(); (7)
gateway.setSerializer(); (8)
gateway.setReceiveTimeout(); (9)
gateway.setExtractPayload(); (10)
gateway.setRecoveryInterval(); (11)
return gateway;
}
<int-redis:queue-inbound-gateway
request-channel="" (1)
reply-channel="" (2)
executor="" (3)
reply-timeout="" (4)
connection-factory="" (5)
queue="" (6)
phase="" (7)
serializer="" (8)
receive-timeout="" (9)
extract-payload="" (10)
recovery-interval=""/> (11)
| 1 | このエンドポイントが Redis データから作成された Message インスタンスを送信する MessageChannel。 |
| 2 | このエンドポイントが応答 Message インスタンスを待機する MessageChannel。オプション - replyChannel ヘッダーはまだ使用されています。 |
| 3 | Spring TaskExecutor (または標準の JDK Executor)Bean への参照。基礎となるリスニングタスクに使用されます。デフォルトは SimpleAsyncTaskExecutor です。 |
| 4 | 応答メッセージが送信されるまで待機するタイムアウト(ミリ秒単位)。通常、キューベースの制限された応答チャネルに適用されます。 |
| 5 | RedisConnectionFactory Bean への参照。デフォルトは redisConnectionFactory です。redis-template 属性とは排他的です。 |
| 6 | 会話 UUID の Redis リストの名前。 |
| 7 | このエンドポイントが開始されるフェーズを指定する SmartLifecycle 属性。デフォルトは 0 です。 |
| 8 | RedisSerializer Bean リファレンス。これは、「シリアライザーなし」を意味する空の文字列にすることができます。この場合、受信 Redis メッセージからの生の byte[] が Message ペイロードとして channel に送信されます。デフォルトは JdkSerializationRedisSerializer です。(バージョン 4.3 より前のリリースでは、デフォルトで StringRedisSerializer であったことに注意してください。その動作を復元するには、StringRedisSerializer への参照を提供します)。 |
| 9 | 受信メッセージが取得されるまでのタイムアウト(ミリ秒単位)。これは通常、キューベースの制限付きリクエストチャネルに適用されます。 |
| 10 | このエンドポイントがペイロードのみを処理するか、Redis キューに対するメッセージ全体を処理するかを指定します。デフォルトは true です。 |
| 11 | 「右ポップ」操作の例外の後、リスナタスクを再起動する前にリスナタスクがスリープする時間(ミリ秒単位)。 |
task-executor は複数のスレッドで処理するように構成する必要があります。そうしないと、エラー発生後に RedisQueueMessageDrivenEndpoint がリスナータスクの再起動を試みたときにデッドロックが発生する可能性があります。errorChannel を使用してこれらのエラーを処理し、再起動を回避することは可能ですが、アプリケーションがデッドロック状態に陥らないようにすることが推奨されます。TaskExecutor の実装については、Spring Framework 参考マニュアルを参照してください。 |
Redis ストリーム送信チャネルアダプター
Spring Integration 5.4 は、メッセージペイロードを Redis ストリームに書き込むためのリアクティブ Redis ストリーム送信チャネルアダプターを導入しました。送信チャネルアダプターは、ReactiveStreamOperations.add(…) を使用して Record をストリームに追加します。次の例は、Redis ストリーム送信チャネルアダプターの Java 構成とサービスクラスを使用する方法を示しています。
Java DSL
Java
@Bean
IntegrationFlow streamOutboundChannelAdapterFlow(ReactiveRedisConnectionFactory connectionFactory) {
return flow -> flow.handle(Redis
.streamOutboundChannelAdapter(connectionFactory, "STREAM_KEY") (1)
.serializationContext() (2)
.hashMapper() (3)
.extractPayload() (4)
);
}
@Bean
@ServiceActivator(inputChannel = "messageChannel")
public ReactiveRedisStreamMessageHandler reactiveValidatorMessageHandler(
ReactiveRedisConnectionFactory reactiveRedisConnectionFactory) {
ReactiveRedisStreamMessageHandler reactiveStreamMessageHandler =
new ReactiveRedisStreamMessageHandler(reactiveRedisConnectionFactory, "myStreamKey"); (1)
reactiveStreamMessageHandler.setSerializationContext(serializationContext); (2)
reactiveStreamMessageHandler.setHashMapper(hashMapper); (3)
reactiveStreamMessageHandler.setExtractPayload(true); (4)
return reactiveStreamMessageHandler;
}
| 1 | ReactiveRedisConnectionFactory とストリーム名を使用して ReactiveRedisStreamMessageHandler のインスタンスを作成し、レコードを追加します。別のコンストラクターバリアントは、SpEL 式に基づいて、リクエストメッセージに対してストリームキーを評価します。 |
| 2 | ストリームに追加する前に、レコードのキーと値を直列化するために使用する RedisSerializationContext を設定します。 |
| 3 | Java 型と Redis ハッシュ / マップ間の契約を提供する HashMapper を設定します。 |
| 4 | 'true' の場合、チャネルアダプターはストリームレコードを追加するリクエストメッセージからペイロードを抽出します。または、メッセージ全体を値として使用します。デフォルトは true です。 |
バージョン 6.5 以降、ReactiveRedisStreamMessageHandler は、内部 ReactiveStreamOperations.add(Record<K, ?> record, XAddOptions xAddOptions) 呼び出しのリクエストメッセージに基づいて RedisStreamCommands.XAddOptions を構築するための setAddOptionsFunction(Function<Message<?>, RedisStreamCommands.XAddOptions> addOptionsFunction) を提供します。
Redis ストリーム受信チャネルアダプター
Spring Integration 5.4 では、Redis ストリームからメッセージを読み取るための Reactive Stream 受信チャネルアダプターが導入されました。受信チャネルアダプターは、自動確認フラグに基づいて StreamReceiver.receive(…) または StreamReceiver.receiveAutoAck() を使用し、Redis ストリームからレコードを読み取ります。次の例は、Redis ストリーム受信チャネルアダプターで Java 構成を使用する方法を示しています。
Java DSL
Java
@Bean
IntegrationFlow streamInboundChannelAdapterFlow(ReactiveRedisConnectionFactory connectionFactory) {
return IntegrationFlow.from(Redis
.streamInboundChannelAdapter(connectionFactory, "STREAM_KEY") (1)
.streamReceiverOptions( (2)
StreamReceiver.StreamReceiverOptions.builder()
.pollTimeout(Duration.ofMillis(100))
.build())
.autoStartup() (3)
.autoAck() (4)
.createConsumerGroup() (5)
.consumerGroup() (6)
.consumerName() (7)
.readOffset() (9)
.extractPayload() (10)
)
.channel(c -> c.flux()) (8)
.get();
}
@Bean
public ReactiveRedisStreamMessageProducer reactiveRedisStreamProducer(
ReactiveRedisConnectionFactory reactiveRedisConnectionFactory) {
ReactiveRedisStreamMessageProducer messageProducer =
new ReactiveRedisStreamMessageProducer(reactiveRedisConnectionFactory, "myStreamKey"); (1)
messageProducer.setStreamReceiverOptions( (2)
StreamReceiver.StreamReceiverOptions.builder()
.pollTimeout(Duration.ofMillis(100))
.build());
messageProducer.setAutoStartup(true); (3)
messageProducer.setAutoAck(false); (4)
messageProducer.setCreateConsumerGroup(true); (5)
messageProducer.setConsumerGroup("my-group"); (6)
messageProducer.setConsumerName("my-consumer"); (7)
messageProducer.setOutputChannel(fromRedisStreamChannel); (8)
messageProducer.setReadOffset(ReadOffset.latest()); (9)
messageProducer.extractPayload(true); (10)
return messageProducer;
}
| 1 | ReactiveRedisConnectionFactory とストリームキーを使用して ReactiveRedisStreamMessageProducer のインスタンスを作成し、レコードを読み取ります。 |
| 2 | リアクティブインフラストラクチャを使用して redis ストリームを消費する StreamReceiver.StreamReceiverOptions。 |
| 3 | アプリケーションコンテキストの起動後にこのエンドポイントを自動的に起動するかどうかを指定するための SmartLifecycle 属性。デフォルトは true です。false の場合、RedisStreamMessageProducer は手動で起動する必要があります(messageProducer.start())。 |
| 4 | false の場合、受信メッセージは自動確認されません。メッセージの確認応答は、クライアントがメッセージを消費するまで延期されます。デフォルトは true です。 |
| 5 | true の場合、コンシューマーグループが作成されます。コンシューマーグループの作成時に、ストリームも作成されます(まだ存在しない場合)。コンシューマーグループはメッセージの配信を追跡し、コンシューマーを区別します。デフォルトは false です。 |
| 6 | コンシューマーグループ名を設定します。デフォルトは定義済みの Bean 名です。 |
| 7 | コンシューマー名を設定します。グループ my-group からメッセージを my-consumer として読み取ります。 |
| 8 | このエンドポイントからメッセージを送信するメッセージチャネル。 |
| 9 | メッセージを読み取るオフセットを定義します。デフォルトは ReadOffset.latest() です。 |
| 10 | "true" の場合、チャネルアダプターは Record からペイロード値を抽出します。それ以外の場合、Record 全体がペイロードとして使用されます。デフォルトは true です。 |
autoAck が false に設定されている場合、Redis ストリームの Record は Redis ドライバーによって自動的に確認応答されず、代わりに IntegrationMessageHeaderAccessor.ACKNOWLEDGMENT_CALLBACK ヘッダーがメッセージに追加され、SimpleAcknowledgment インスタンスが値として生成されます。このようなレコードに基づいてメッセージのビジネスロジックが実行されるたびに、acknowledge() コールバックを呼び出すのは、ターゲット統合フローの責任です。デシリアライズ中に例外が発生し、errorChannel が設定されている場合にも、同様のロジックが必要です。そのため、ターゲットエラーハンドラーは、このような失敗したメッセージを確認するか、確認応答しないかを決定する必要があります。IntegrationMessageHeaderAccessor.ACKNOWLEDGMENT_CALLBACK に加えて、ReactiveRedisStreamMessageProducer もこれらのヘッダーをメッセージに設定して、RedisHeaders.STREAM_KEY、RedisHeaders.STREAM_MESSAGE_ID、RedisHeaders.CONSUMER_GROUP、RedisHeaders.CONSUMER を生成します。
バージョン 5.5 以降、StreamReceiver.StreamReceiverOptionsBuilder オプションを ReactiveRedisStreamMessageProducer で明示的に設定できるようになりました。これには、新たに導入された onErrorResume 関数も含まれます。この関数は、デシリアライゼーションエラーが発生した際に Redis ストリームコンシューマーがポーリングを継続する必要がある場合に必要です。デフォルトの関数は、前述のように、失敗したメッセージの確認応答を含むメッセージをエラーチャネル(提供されている場合)に送信します。これらの StreamReceiver.StreamReceiverOptionsBuilder はすべて、外部から提供される StreamReceiver.StreamReceiverOptions と相互に排他的です。
Redis ロックレジストリ
Spring Integration 4.0 では RedisLockRegistry が導入されました。特定のコンポーネント(アグリゲータやリシーケンサなど)は、LockRegistry インスタンスから取得したロックを使用して、一度に 1 つのスレッドのみがグループを操作できるようにします。DefaultLockRegistry は、この機能を単一のコンポーネント内で実行します。これらのコンポーネントには、外部ロックレジストリを設定できます。共有 MessageGroupStore と併用することで、RedisLockRegistry は複数のアプリケーションインスタンスにこの機能を提供するように設定でき、一度に 1 つのインスタンスのみがグループを操作できるようになります。
ローカルスレッドによってロックが解除されると、通常、別のローカルスレッドがすぐにロックを取得できます。別のレジストリインスタンスを使用するスレッドによってロックが解放された場合、ロックを取得するのに最大 100 ミリ秒かかることがあります。
サーバーの障害発生時にロックが「ハング」するのを防ぐため、このレジストリ内のロックはデフォルトで 60 秒後に期限切れとなりますが、レジストリで設定することも可能です。ロックは通常、それよりもずっと短い時間保持されます。
| キーは有効期限切れになる可能性があるため、有効期限切れのロックを解除しようとすると例外がスローされます。ただし、そのようなロックによって保護されているリソースが侵害されている可能性があるため、このような例外は重大な問題と見なす必要があります。有効期限は、このような状況を防ぐのに十分な大きさに設定する必要がありますが、サーバー障害発生後にロックを適切な時間内に回復できる程度に低く設定する必要があります。 |
バージョン 5.0 から、RedisLockRegistry は ExpirableLockRegistry を実装します。ExpirableLockRegistry は、age より前に最後に取得され、現在ロックされていないロックを削除します。
バージョン 5.5.6 以降、RedisLockRegistry は、RedisLockRegistry.setCacheCapacity() を介して未使用の RedisLock インスタンスの自動キャッシュクリーンアップをサポートします。デフォルト値は 100_000 です。詳細については、JavaDocs を参照してください。キャッシュは、expireUnusedOlderThan(long age) API を介してもクリーンアップできます。キャッシュされたすべての RedisLock インスタンスが取得され、それらのいずれも削除できない場合、obtain(Object lockKey) メソッドは CannotAcquireLockException をスローします。キャッシュされていない RedisLock をロックしようとすると、別の CannotAcquireLockException もスローされます。そのロックキーに対して新しい obtain(Object lockKey) を呼び出す必要があります。
バージョン 5.5.13 以降、RedisLockRegistry は、Redis ロックの取得をどのモードで行うかを決定する setRedisLockType(RedisLockType) オプションを公開します。
RedisLockType.SPIN_LOCK- ロックは、ロックを取得できるかどうかをチェックする定期的なループ(100ms)によって取得されます。デフォルト。RedisLockType.PUB_SUB_LOCK- ロックは、redispub-sub サブスクリプションによって取得されます。
pub-sub モードは推奨モードです。クライアント Redis サーバー間のネットワーク通信が少なく、パフォーマンスも向上します。他のプロセスでサブスクリプションのロック解除が通知されると、すぐにロックが取得されます。ただし、Redis はマスター / レプリカ接続(AWS ElastiCache 環境など)で pub-sub モードをサポートしていないため、レジストリをあらゆる環境で動作させるため、デフォルトで busy-spin モードが選択されています。
バージョン 6.4 以降では、ロックの所有権が期限切れの場合、RedisLockRegistry.RedisLock.unlock() メソッドは IllegalStateException をスローする代わりに ConcurrentModificationException をスローします。
バージョン 6.4 以降では、ロックの定期的な更新のスケジューラを構成するための RedisLockRegistry.setRenewalTaskScheduler() が追加されました。これを設定すると、ロックが正常に取得されてから、ロックが解除されるか、Redis キーが削除されるまで、有効期限の 1/3 ごとにロックが自動的に更新されます。
バージョン 7.0 以降、RedisLock は DistributedLock インターフェースを実装し、ロックステータスデータの TTL(Time To Live)をカスタマイズする機能をサポートします。RedisLock は、lock(Duration ttl) または tryLock(long time, TimeUnit unit, Duration ttl) メソッドを使用して、指定された TTL 値で取得できるようになりました。RedisLockRegistry では、新しい renewLock(Object lockKey, Duration ttl) メソッドが提供され、カスタム TTL 値でロックを更新できるようになりました。
7.1 以降のバージョンでは、RedisLockRegistry は Redis 8.4 以降のバージョンに対して、ネイティブの Redis CAS(比較設定)および CAD(比較削除)コマンドを使用します。それ以前の Redis バージョンでは、レジストリは自動的に以前の Lua スクリプトベースのアプローチにフォールバックします。
クラスターモードでの Valkey サポート用の AWS ElastiCache
バージョン 6.4.9/6.5.4/7.0.0 以降、RedisLockRegistry はクラスターモードで AWS Elasticache for Valkey をサポートします。このバージョンの Valkey(Redis の代替)では、すべての PubSub 操作(PUBLISH、SUBSCRIBE など)は、内部的にシャード化されたバリアント(SPUBLISH、SSUBSCRIBE など)を使用します。以下の形式のエラーが発生した場合:
Caused by: io.lettuce.core.RedisCommandExecutionException: ERR Script attempted to access keys that do not hash to the same slot script: b2dedc0ab01c17f9f20e3e6ddb62dcb6afbed0bd, on @user_script:3.RedisLockRegistry の unlock ステップでは、ハッシュタグ {…} を含むロックキーを指定して、unlock スクリプト内のすべての操作が同じクラスタースロット / シャードにハッシュされるようにする必要があります。例:
RedisLockRegistry lockRegistry = new RedisLockRegistry("my-lock-key{choose_your_tag}");
lockRegistry.lock();
# critical section
lockRegistry.unlock();