インターフェースの使用
org.springframework.kafka.core.KafkaOperations
KafkaOperations を使用するパッケージ
パッケージ
説明
kafka コアコンポーネントのパッケージ
kafka リスナーのためのパッケージ
リクエスト / 応答セマンティクスのクラスを提供します。
再試行可能なトピック処理用のパッケージ。
org.springframework.kafka.core 内の KafkaOperations 使用
KafkaOperations を実装している org.springframework.kafka.core のクラス修飾子と型クラス説明classKafkaTemplate<K,V> 高レベルの操作を実行するためのテンプレート。classトピック名に基づいてメッセージをルーティングするKafkaTemplate。型 KafkaOperations のパラメーターを持つ org.springframework.kafka.core のメソッド修飾子と型メソッド説明@Nullable TKafkaOperations.OperationsCallback.doInOperations(KafkaOperations<K, V> operations) org.springframework.kafka.listener 内の KafkaOperations 使用
型 KafkaOperations のパラメーターを持つ org.springframework.kafka.listener のメソッド修飾子と型メソッド説明protected DurationSEDeadLetterPublishingRecoverer.determineSendTimeout(KafkaOperations<?, ?> template) テンプレートのプロデューサーファクトリとDeadLetterPublishingRecoverer.setWaitForSendResultTimeout(Duration)に基づいて送信タイムアウトを決定します。protected voidDeadLetterPublishingRecoverer.publish(org.apache.kafka.clients.producer.ProducerRecord<ObjectSE, ObjectSE> outRecord, KafkaOperations<ObjectSE, ObjectSE> kafkaTemplate, org.apache.kafka.clients.consumer.ConsumerRecord<?, ?> inRecord) sendの結果のログ記録以上のものが必要な場合は、これをオーバーライドします。protected voidDeadLetterPublishingRecoverer.send(org.apache.kafka.clients.producer.ProducerRecord<ObjectSE, ObjectSE> outRecord, KafkaOperations<ObjectSE, ObjectSE> kafkaTemplate, org.apache.kafka.clients.consumer.ConsumerRecord<?, ?> inRecord) レコードを送信します。protected voidDeadLetterPublishingRecoverer.verifySendResult(KafkaOperations<ObjectSE, ObjectSE> kafkaTemplate, org.apache.kafka.clients.producer.ProducerRecord<ObjectSE, ObjectSE> outRecord, @Nullable CompletableFutureSE<SendResult<ObjectSE, ObjectSE>> sendResult, org.apache.kafka.clients.consumer.ConsumerRecord<?, ?> inRecord) sendの将来が完了するまで待ちます。型 KafkaOperations のパラメーターを持つ org.springframework.kafka.listener のコンストラクター修飾子コンストラクター説明DeadLetterPublishingRecoverer(KafkaOperations<?, ?> template) 提供されたテンプレートと、失敗したレコードの元のトピック ( "-dlt" が追加された) と失敗したレコードと同じパーティションに基づいて TopicPartition を返すデフォルトの宛先解決関数を使用してインスタンスを作成します。DeadLetterPublishingRecoverer(KafkaOperations<?, ?> template, BiFunctionSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE, @Nullable org.apache.kafka.common.TopicPartition> destinationResolver) 提供されたテンプレートと宛先解決関数を使用してインスタンスを作成し、失敗したコンシューマーレコードと例外を受け取ってTopicPartitionを返します。DefaultAfterRollbackProcessor(@Nullable BiConsumerSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE> recoverer, BackOff backOff, @Nullable KafkaOperations<?, ?> kafkaOperations, boolean commitRecovered) backOff がトピック / パーティション / オフセットに対して STOP を返した後に呼び出される、提供されたリカバリを使用してインスタンスを構築します。DefaultAfterRollbackProcessor(@Nullable BiConsumerSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE> recoverer, BackOff backOff, @Nullable BackOffHandler backOffHandler, @Nullable KafkaOperations<?, ?> kafkaOperations, boolean commitRecovered) backOff がトピック / パーティション / オフセットに対して STOP を返した後に呼び出される、提供されたリカバリを使用してインスタンスを構築します。型の型引数を持つ org.springframework.kafka.listener のコンストラクターパラメーター KafkaOperations修飾子コンストラクター説明DeadLetterPublishingRecoverer(FunctionSE<org.apache.kafka.clients.producer.ProducerRecord<?, ?>, ? extends @Nullable KafkaOperations<?, ?>> templateResolver, boolean transactional, BiFunctionSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE, @Nullable org.apache.kafka.common.TopicPartition> destinationResolver) 失敗したコンシューマーレコードと例外を受け取り、KafkaOperationsと、このインスタンスからの公開がトランザクションかどうかのフラグを返すテンプレート解決関数を使用してインスタンスを作成します。DeadLetterPublishingRecoverer(FunctionSE<org.apache.kafka.clients.producer.ProducerRecord<?, ?>, ? extends @Nullable KafkaOperations<?, ?>> templateResolver, BiFunctionSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE, @Nullable org.apache.kafka.common.TopicPartition> destinationResolver) 失敗したコンシューマーレコードと例外を受け取り、KafkaOperationsと、このインスタンスからの公開がトランザクションかどうかのフラグを返すテンプレート解決関数を使用してインスタンスを作成します。DeadLetterPublishingRecoverer(MapSE<ClassSE<?>, ? extends @Nullable KafkaOperations<?, ?>> templates) 提供されたテンプレートと、失敗したレコードの元のトピック ( "-dlt" が追加された) と失敗したレコードと同じパーティションに基づいて TopicPartition を返すデフォルトの宛先解決関数を使用してインスタンスを作成します。DeadLetterPublishingRecoverer(MapSE<ClassSE<?>, ? extends @Nullable KafkaOperations<?, ?>> templates, BiFunctionSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE, @Nullable org.apache.kafka.common.TopicPartition> destinationResolver) 提供されたテンプレートと宛先解決関数を使用して、失敗したコンシューマーレコードと例外を受け取り、TopicPartitionを返すインスタンスを作成します。org.springframework.kafka.requestreply 内の KafkaOperations 使用
KafkaOperations を実装している org.springframework.kafka.requestreply のクラス修飾子と型クラス説明class同じ相関 ID を持つ複数の応答を集約する応答テンプレート。classReplyingKafkaTemplate<K,V, R> リクエスト / 応答セマンティクスを実装する KafkaTemplate。org.springframework.kafka.retrytopic 内の KafkaOperations 使用
型 KafkaOperations のパラメーターを持つ org.springframework.kafka.retrytopic のメソッド修飾子と型メソッド説明RetryTopicConfigurationBuilder.create(KafkaOperations<?, ?> sendToTopicKafkaTemplate) 提供されたテンプレートを使用してRetryTopicConfigurationを作成します。型 KafkaOperations の型引数を持つ org.springframework.kafka.retrytopic のメソッドパラメーター修飾子と型メソッド説明DeadLetterPublishingRecovererFactory.DeadLetterPublisherCreator.create(FunctionSE<org.apache.kafka.clients.producer.ProducerRecord<?, ?>, ? extends @Nullable KafkaOperations<?, ?>> templateResolver, BiFunctionSE<org.apache.kafka.clients.consumer.ConsumerRecord<?, ?>, ExceptionSE, @Nullable org.apache.kafka.common.TopicPartition> destinationResolver) 提供されたプロパティを使用してDeadLetterPublishingRecovererを作成します。型 KafkaOperations のパラメーターを持つ org.springframework.kafka.retrytopic のコンストラクター修飾子コンストラクター説明DestinationTopicPropertiesFactory(@Nullable StringSE retryTopicSuffix, @Nullable StringSE dltSuffix, ListSE<LongSE> backOffValues, ExceptionMatcher exceptionMatcher, int numPartitions, KafkaOperations<?, ?> kafkaOperations, DltStrategy dltStrategy, TopicSuffixingStrategy topicSuffixingStrategy, SameIntervalTopicReuseStrategy sameIntervalTopicReuseStrategy, long timeout, MapSE<StringSE, SetSE<ClassSE<? extends ThrowableSE>>> dltRoutingRules) 提供されたプロパティを使用してインスタンスを構築します。DestinationTopicPropertiesFactory(StringSE retryTopicSuffix, StringSE dltSuffix, ListSE<LongSE> backOffValues, ExceptionMatcher exceptionMatcher, int numPartitions, KafkaOperations<?, ?> kafkaOperations, DltStrategy dltStrategy, TopicSuffixingStrategy topicSuffixingStrategy, SameIntervalTopicReuseStrategy sameIntervalTopicReuseStrategy, long timeout) 提供されたプロパティを使用してインスタンスを構築します。Properties(long delayMs, StringSE suffix, org.springframework.kafka.retrytopic.DestinationTopic.Type type, int maxAttempts, int numPartitions, DltStrategy dltStrategy, KafkaOperations<?, ?> kafkaOperations, BiPredicateSE<IntegerSE, ThrowableSE> shouldRetryOn, long timeout) DLT コンテナーが自動的に起動するように、提供されたプロパティを使用してインスタンスを作成します(コンテナーファクトリがそのように構成されている場合)。Properties(long delayMs, StringSE suffix, org.springframework.kafka.retrytopic.DestinationTopic.Type type, int maxAttempts, int numPartitions, DltStrategy dltStrategy, KafkaOperations<?, ?> kafkaOperations, BiPredicateSE<IntegerSE, ThrowableSE> shouldRetryOn, long timeout, @Nullable BooleanSE autoStartDltHandler) 提供されたプロパティでインスタンスを作成します。Properties(long delayMs, StringSE suffix, org.springframework.kafka.retrytopic.DestinationTopic.Type type, int maxAttempts, int numPartitions, DltStrategy dltStrategy, KafkaOperations<?, ?> kafkaOperations, BiPredicateSE<IntegerSE, ThrowableSE> shouldRetryOn, long timeout, @Nullable BooleanSE autoStartDltHandler, SetSE<ClassSE<? extends ThrowableSE>> usedForExceptions) 提供されたプロパティでインスタンスを作成します。