クラス ReactiveRedisStreamMessageProducer

実装済みのインターフェース一覧:
Aware, BeanFactoryAware, BeanNameAware, DisposableBean, InitializingBean, SmartInitializingSingleton, ApplicationContextAware, Lifecycle, Phased, SmartLifecycle, ComponentSourceAware, ExpressionCapable, MessageProducer, IntegrationPattern, NamedComponent, IntegrationInboundManagement, IntegrationManagement, ManageableLifecycle, ManageableSmartLifecycle, TrackableComponent

public class ReactiveRedisStreamMessageProducer extends MessageProducerSupport
Redis ストリームからメッセージを読み取り、提供された出力チャネルに公開するための MessageProducerSupport。デフォルトでは、このアダプターはメッセージをスタンドアロンクライアント XREAD (Redis コマンド)として読み取りますが、consumerName フィールドを設定することにより、コンシューマーグループ機能 XREADGROUP に切り替えることができます。デフォルトでは、コンシューマーグループ名はこの Bean IntegrationObjectSupport.getBeanName() の ID です。
導入:
5.4
作成者:
Attoumane Ahamadi, Artem Bilan, Rohan Mukesh
  • コンストラクターの詳細

  • 方法の詳細

    • setReadOffset

      public void setReadOffset(ReadOffset readOffset)
      メッセージを読み取るオフセットを定義します。デフォルトでは、ReadOffset.latest() が使用されます。ReadOffset.latest() は、ストリームに追加された新しいデータを取得するために XREAD で使用される ID である "$" と同じです。コンシューマーグループ機能に切り替えるときに、それが ReadOffset.latest() と等しい場合は、ReadOffset.lastConsumed() に設定することに注意してください。
      パラメーター:
      readOffset - 希望のオフセット
    • setExtractPayload

      public void setExtractPayload(boolean extractPayload)
      このチャネルアダプターを構成して、Record から値を抽出するかどうかを指定します。
      パラメーター:
      extractPayload - デフォルト true
    • setAutoAck

      public void setAutoAck(boolean autoAck)
      コンシューマーグループで読み取ったメッセージを確認するかどうかを設定します。デフォルトでは true
      パラメーター:
      autoAck - 確認オプション。
    • setConsumerGroup

      public void setConsumerGroup(StringSE consumerGroup)
      コンシューマーグループの名前を設定します。必要に応じて、そのコンシューマーグループを作成することができます。createConsumerGroup を参照してください。設定されていない場合、定義された Bean 名 IntegrationObjectSupport.getBeanName() が使用されます。
      パラメーター:
      consumerGroup - このアダプターがメッセージをリッスンするために登録する必要があるコンシューマーグループ。
    • setConsumerName

      public void setConsumerName(@Nullable StringSE consumerName)
      コンシューマーの名前を設定します。コンシューマー名が指定されると、このアダプターはコンシューマーグループ機能に切り替えられます。この値はグループ内で一意である必要があることに注意してください。
      パラメーター:
      consumerName - コンシューマーグループのコンシューマー名
    • setCreateConsumerGroup

      public void setCreateConsumerGroup(boolean createConsumerGroup)
      コンシューマーグループが存在しない場合にのみ、コンシューマーグループを作成します。作成中に、ストリームも作成します。MKSTREAM を参照してください。
      パラメーター:
      createConsumerGroup - コンシューマーグループ、デフォルトで false を作成する必要があるかどうかを指定します
    • setStreamReceiverOptions

      public void setStreamReceiverOptions(@Nullable StreamReceiver.StreamReceiverOptions<StringSE,?> streamReceiverOptions)
      StreamReceiver をカスタマイズするために使用される ReactiveStreamOperations を設定します。ポーリングタイムアウトと直列化コンテキストを設定する方法を提供します。デフォルトでは、ポーリングタイムアウトは無限に設定され、StringRedisSerializer が使用されます。'pollTimeout'、'batchSize'、'onErrorResume'、'serializer'、'targetType'、'objectMapper' とは相互に排他的です。
      パラメーター:
      streamReceiverOptions - 必要なレシーバーオプション
    • setPollTimeout

      public void setPollTimeout(DurationSE pollTimeout)
      読み取り中の BLOCK オプションのポーリングタイムアウトを設定します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。
      パラメーター:
      pollTimeout - ポーリングのタイムアウト。
      導入:
      5.5
      関連事項:
    • setBatchSize

      public void setBatchSize(int recordsPerPoll)
      読み取り中に COUNT オプションのバッチサイズを構成します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。
      パラメーター:
      recordsPerPoll - ゼロより大きくなければなりません。
      導入:
      5.5
      関連事項:
    • setOnErrorResume

      public void setOnErrorResume(FunctionSE<? super ThrowableSE, ? extends org.reactivestreams.Publisher<VoidSE>> resumeFunction)
      ストリームのポーリングが失敗したときにメインシーケンスを再開するように再開機能を構成します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。デフォルトでは、この関数は失敗した Record を抽出し、提供された MessageProducerSupport.setErrorChannel(MessageChannel)ErrorMessage を送信します。このメッセージプロデューサーに手動確認応答が構成されている場合、このレコードの失敗したメッセージには IntegrationMessageHeaderAccessor.ACKNOWLEDGMENT_CALLBACK ヘッダーが含まれる場合があります。
      パラメーター:
      resumeFunction - null であってはなりません。
      導入:
      5.5
      関連事項:
    • setSerializer

      public void setSerializer(RedisSerializationContext.SerializationPair<?> pair)
      キー、ハッシュキー、ハッシュ値シリアライザーを構成します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。
      パラメーター:
      pair - null であってはなりません。
      導入:
      5.5
      関連事項:
    • setTargetType

      public void setTargetType(ClassSE<?> targetType)
      ハッシュターゲット型を構成します。発行されたレコード型を ObjectRecord に変更します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。
      パラメーター:
      targetType - null であってはなりません。
      導入:
      5.5
      関連事項:
    • setObjectMapper

      public void setObjectMapper(HashMapper<?,?,?> hashMapper)
      ハッシュマッパーを構成します。setStreamReceiverOptions(StreamReceiver.StreamReceiverOptions) と相互に排他的です。
      パラメーター:
      hashMapper - null であってはなりません。
      導入:
      5.5
      関連事項:
    • getComponentType

      public StringSE getComponentType()
      次で指定:
      インターフェース NamedComponent 内の getComponentType 
      オーバーライド:
      クラス MessageProducerSupportgetComponentType 
    • onInit

      protected void onInit()
      クラスからコピーされた説明: IntegrationObjectSupport
      サブクラスは、初期化ロジック用にこれを実装できます。
      オーバーライド:
      クラス MessageProducerSupportonInit 
    • doStart

      protected void doStart()
      クラスからコピーされた説明: MessageProducerSupport
      デフォルトでは何も実行されません。ライフサイクル管理された動作が必要な場合、サブクラスはこれをオーバーライドできます。'lifecycleLock' によって保護されています。
      オーバーライド:
      クラス MessageProducerSupportdoStart