クラス KafkaAdmin
java.lang.ObjectSE
org.springframework.kafka.core.KafkaResourceFactory
org.springframework.kafka.core.KafkaAdmin
- 実装済みのインターフェース一覧:
Aware, SmartInitializingSingleton, ApplicationContextAware, KafkaAdminOperations
public class KafkaAdmin
extends KafkaResourceFactory
implements ApplicationContextAware, SmartInitializingSingleton, KafkaAdminOperations
アプリケーションコンテキストで定義されたトピックを作成するために
Admin に委譲する管理者。- 導入:
- 1.3
- 作成者:
- Gary Russell, Artem Bilan, Adrian Gygax, Sanghyeok An, Valentina Armenise, Anders Swanson, Omer Celik, Choi Wang Gyu, Go Beom Jun
ネストされたクラスの概要
ネストされたクラス修飾子と型クラス説明static class複数のトピックを単一の Bean として宣言するのを容易にするNewTopicのコレクションのラッパー。フィールド概要
フィールド修飾子と型フィールド説明static final DurationSEデフォルトのクローズタイムアウト期間は 10 秒です。static final DurationSEThe default interval between attempts to obtain the cluster id after a failed attempt, as 5 minutes.コンストラクター概要
コンストラクター方法の概要
修飾子と型メソッド説明void@Nullable StringSEReturn the cluster id, fetching it from the broker on the first call.protected org.apache.kafka.clients.admin.AdminAdminClientクラスを使用して新しいAdminクライアントインスタンスを作成します。voidcreateOrModifyTopics(org.apache.kafka.clients.admin.NewTopic... topics) トピックが存在しない場合は作成するか、必要に応じてパーティションの数を増やします。voiddeleteTopics(StringSE... topicNames) Kafka クラスターからトピックを削除します。describeTopics(StringSE... topicNames) これらのトピックのTopicDescriptionを取得します。@Nullable StringSEclusterId プロパティを取得します。この管理者の構成の変更不可能なコピーを取得します。protected PredicateSE<org.apache.kafka.clients.admin.NewTopic> NewTopicを作成または変更するかどうかを決定するために使用される述語を返します。int操作のタイムアウトを秒単位で返します。final booleanこのメソッドを呼び出して、トピックをチェック / 追加します。これは、アプリケーションコンテキストが初期化されたときにブローカーが使用できず、fatalIfBrokerNotAvailableが false であるか、autoCreateが false に設定されている場合に必要になることがあります。protected CollectionSE<org.apache.kafka.clients.admin.NewTopic> 作成または変更するNewTopicのコレクションを返します。voidsetApplicationContext(ApplicationContext applicationContext) voidsetAutoCreate(boolean autoCreate) コンテキストの初期化中にトピックの自動作成を抑制するには、false に設定します。voidsetCloseTimeout(int closeTimeout) クローズタイムアウトを秒単位で設定します。voidsetClusterId(StringSE clusterId) クラスター ID を設定します。voidsetClusterIdRetryInterval(DurationSE clusterIdRetryInterval) Set the minimum interval between attempts to fetch the cluster id from the broker after a failed attempt.voidsetCreateOrModifyTopic(PredicateSE<org.apache.kafka.clients.admin.NewTopic> createOrModifyTopic) 検出されたNewTopicBean がこの管理インスタンスによる作成または変更の対象となる場合に true を返す述語を設定します。voidsetFatalIfBrokerNotAvailable(boolean fatalIfBrokerNotAvailable) 初期化中にブローカーに接続できない場合にアプリケーションコンテキストのロードに失敗する場合は、true に設定して、トピックを確認 / 追加します。voidsetModifyTopicConfigs(boolean modifyTopicConfigs) true に設定すると、現在のトピック構成プロパティがNewTopicBean のプロパティと比較され、異なる場合は更新されます。voidsetOperationTimeout(int operationTimeout) 操作タイムアウトを秒単位で設定します。クラス KafkaResourceFactory から継承されたメソッド
checkBootstrap, getBootstrapServers, setBootstrapServersSupplier
フィールドの詳細
DEFAULT_CLOSE_TIMEOUT
デフォルトのクローズタイムアウト期間は 10 秒です。DEFAULT_CLUSTER_ID_RETRY_INTERVAL
The default interval between attempts to obtain the cluster id after a failed attempt, as 5 minutes.- 導入:
- 4.1.1
コンストラクターの詳細
KafkaAdmin
方法の詳細
setApplicationContext
- 次で指定:
- インターフェース
ApplicationContextAware内のsetApplicationContext - 例外:
BeansException
setCloseTimeout
public void setCloseTimeout(int closeTimeout) クローズタイムアウトを秒単位で設定します。デフォルトはDEFAULT_CLOSE_TIMEOUT秒です。- パラメーター:
closeTimeout- タイムアウト。
setOperationTimeout
public void setOperationTimeout(int operationTimeout) 操作のタイムアウトを秒単位で設定します。デフォルトは 30 秒です。- パラメーター:
operationTimeout- タイムアウト。
getOperationTimeout
public int getOperationTimeout()操作のタイムアウトを秒単位で返します。- 戻り値:
- タイムアウト。
- 導入:
- 3.0.2
setFatalIfBrokerNotAvailable
public void setFatalIfBrokerNotAvailable(boolean fatalIfBrokerNotAvailable) 初期化中にブローカーに接続できない場合にアプリケーションコンテキストのロードに失敗する場合は、true に設定して、トピックを確認 / 追加します。- パラメーター:
fatalIfBrokerNotAvailable- 失敗するのは本当です。
setAutoCreate
public void setAutoCreate(boolean autoCreate) コンテキストの初期化中にトピックの自動作成を抑制するには、false に設定します。- パラメーター:
autoCreate- コンテキストの初期化中にトピックを作成するかどうかを示すブールフラグ- 関連事項:
setModifyTopicConfigs
public void setModifyTopicConfigs(boolean modifyTopicConfigs) true に設定すると、現在のトピック構成プロパティがNewTopicBean のプロパティと比較され、異なる場合は更新されます。- パラメーター:
modifyTopicConfigs- 必要に応じて構成を確認および更新する場合は true。- 導入:
- 2.8.7
setCreateOrModifyTopic
public void setCreateOrModifyTopic(PredicateSE<org.apache.kafka.clients.admin.NewTopic> createOrModifyTopic) 検出されたNewTopicBean がこの管理インスタンスによる作成または変更の対象となる場合に true を返す述語を設定します。デフォルトの述語は、すべてのNewTopicに対して true を返します。newTopics()のデフォルト実装によって使用されます。- パラメーター:
createOrModifyTopic- 述語。- 導入:
- 2.9.10
- 関連事項:
getCreateOrModifyTopic
NewTopicを作成または変更するかどうかを決定するために使用される述語を返します。- 戻り値:
- 述語。
- 導入:
- 2.9.10
- 関連事項:
setClusterId
クラスター ID を設定します。これを使用すると、ユーザーが管理者権限を持っていない場合に、ブローカーからクラスター ID を取得しようとするのを防ぐことができます。- パラメーター:
clusterId- 設定する clusterId- 導入:
- 3.1
getClusterId
setClusterIdRetryInterval
Set the minimum interval between attempts to fetch the cluster id from the broker after a failed attempt. Defaults to5 minutes. Set toDuration.ZEROSE to attempt a fetch on everyclusterId()call.- パラメーター:
clusterIdRetryInterval- 間隔。- 導入:
- 4.1.1
- 関連事項:
getConfigurationProperties
インターフェースからコピーされた説明:KafkaAdminOperationsこの管理者の構成の変更不可能なコピーを取得します。- 次で指定:
- インターフェース
KafkaAdminOperations内のgetConfigurationProperties - 戻り値:
- 構成マップ。
afterSingletonsInstantiated
public void afterSingletonsInstantiated()- 次で指定:
- インターフェース
SmartInitializingSingleton内のafterSingletonsInstantiated
initialize
public final boolean initialize()このメソッドを呼び出して、トピックをチェック / 追加します。これは、アプリケーションコンテキストが初期化されたときにブローカーが使用できず、fatalIfBrokerNotAvailableが false であるか、autoCreateが false に設定されている場合に必要になることがあります。- 戻り値:
- 成功した場合は true。
- 関連事項:
newTopics
作成または変更するNewTopicのコレクションを返します。デフォルトの実装では、アプリケーションコンテキスト内のすべてのNewTopicBean を取得し、それぞれにsetCreateOrModifyTopic(Predicate)述語を適用します。同じトピック名の NewTopic がある場合は、TopicForRetryableBean も削除されます。これは述語を呼び出す前に実行されます。- 戻り値:
NewTopicのコレクション。- 導入:
- 2.9.10
- 関連事項:
clusterId
Return the cluster id, fetching it from the broker on the first call. A failed fetch is remembered and not attempted again until theclusterIdRetryIntervalhas elapsed; each attempt creates anAdmininstance and blocks for up tooperationTimeout, so retrying on every call can stall callers such as an observation-enabled listener container.- 次で指定:
- インターフェース
KafkaAdminOperations内のclusterId - 戻り値:
- the cluster id, or null if it has not been fetched successfully.
createOrModifyTopics
public void createOrModifyTopics(org.apache.kafka.clients.admin.NewTopic... topics) インターフェースからコピーされた説明:KafkaAdminOperationsトピックが存在しない場合は作成するか、必要に応じてパーティションの数を増やします。- 次で指定:
- インターフェース
KafkaAdminOperations内のcreateOrModifyTopics - パラメーター:
topics- トピック。
describeTopics
public MapSE<StringSE, org.apache.kafka.clients.admin.TopicDescription> describeTopics(StringSE... topicNames) インターフェースからコピーされた説明:KafkaAdminOperationsこれらのトピックのTopicDescriptionを取得します。- 次で指定:
- インターフェース
KafkaAdminOperations内のdescribeTopics - パラメーター:
topicNames- トピック名。- 戻り値:
- 名前: topicDescription の地図。
deleteTopics
Kafka クラスターからトピックを削除します。- 次で指定:
- インターフェース
KafkaAdminOperations内のdeleteTopics - パラメーター:
topicNames- 削除するトピック名。- 例外:
KafkaException- 操作が失敗した場合。- 導入:
- 4.0
createAdmin
protected org.apache.kafka.clients.admin.Admin createAdmin()AdminClientクラスを使用して新しいAdminクライアントインスタンスを作成します。- 戻り値:
- 新しい
Adminクライアントインスタンス。 - 導入:
- 3.3.0
- 関連事項:
getAdminConfig