クラス KafkaAdmin

実装済みのインターフェース一覧:
Aware, SmartInitializingSingleton, ApplicationContextAware, 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 DurationSE
    The default interval between attempts to obtain the cluster id after a failed attempt, as 5 minutes.
  • コンストラクター概要

    コンストラクター
    コンストラクター
    説明
    提供された構成に基づいて、Admin を使用してインスタンスを作成します。
  • 方法の概要

    修飾子と型
    メソッド
    説明
    void
    @Nullable StringSE
    Return the cluster id, fetching it from the broker on the first call.
    protected org.apache.kafka.clients.admin.Admin
    AdminClient クラスを使用して新しい Admin クライアントインスタンスを作成します。
    void
    createOrModifyTopics(org.apache.kafka.clients.admin.NewTopic... topics)
    トピックが存在しない場合は作成するか、必要に応じてパーティションの数を増やします。
    void
    deleteTopics(StringSE... topicNames)
    Kafka クラスターからトピックを削除します。
    MapSE<StringSE, org.apache.kafka.clients.admin.TopicDescription>
    describeTopics(StringSE... topicNames)
    これらのトピックの TopicDescription を取得します。
    @Nullable StringSE
    clusterId プロパティを取得します。
    この管理者の構成の変更不可能なコピーを取得します。
    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 のコレクションを返します。
    void
    void
    setAutoCreate(boolean autoCreate)
    コンテキストの初期化中にトピックの自動作成を抑制するには、false に設定します。
    void
    setCloseTimeout(int closeTimeout)
    クローズタイムアウトを秒単位で設定します。
    void
    クラスター ID を設定します。
    void
    setClusterIdRetryInterval(DurationSE clusterIdRetryInterval)
    Set the minimum interval between attempts to fetch the cluster id from the broker after a failed attempt.
    void
    setCreateOrModifyTopic(PredicateSE<org.apache.kafka.clients.admin.NewTopic> createOrModifyTopic)
    検出された NewTopic Bean がこの管理インスタンスによる作成または変更の対象となる場合に true を返す述語を設定します。
    void
    setFatalIfBrokerNotAvailable(boolean fatalIfBrokerNotAvailable)
    初期化中にブローカーに接続できない場合にアプリケーションコンテキストのロードに失敗する場合は、true に設定して、トピックを確認 / 追加します。
    void
    setModifyTopicConfigs(boolean modifyTopicConfigs)
    true に設定すると、現在のトピック構成プロパティが NewTopic Bean のプロパティと比較され、異なる場合は更新されます。
    void
    setOperationTimeout(int operationTimeout)
    操作タイムアウトを秒単位で設定します。

    クラス ObjectSE から継承されたメソッド

    clone, equalsSE, finalize, getClass, hashCode, notify, notifyAll, toString, wait, waitSE, waitSE
  • フィールドの詳細

    • DEFAULT_CLOSE_TIMEOUT

      public static final DurationSE DEFAULT_CLOSE_TIMEOUT
      デフォルトのクローズタイムアウト期間は 10 秒です。
    • DEFAULT_CLUSTER_ID_RETRY_INTERVAL

      public static final DurationSE 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

      public KafkaAdmin(MapSE<StringSE,ObjectSE> config)
      提供された構成に基づいて、Admin を使用してインスタンスを作成します。
      パラメーター:
      config - Admin の構成。
  • 方法の詳細

    • setApplicationContext

      public void setApplicationContext(ApplicationContext applicationContext) throws BeansException
      次で指定:
      インターフェース 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 に設定すると、現在のトピック構成プロパティが NewTopic Bean のプロパティと比較され、異なる場合は更新されます。
      パラメーター:
      modifyTopicConfigs - 必要に応じて構成を確認および更新する場合は true。
      導入:
      2.8.7
    • setCreateOrModifyTopic

      public void setCreateOrModifyTopic(PredicateSE<org.apache.kafka.clients.admin.NewTopic> createOrModifyTopic)
      検出された NewTopic Bean がこの管理インスタンスによる作成または変更の対象となる場合に true を返す述語を設定します。デフォルトの述語は、すべての NewTopic に対して true を返します。newTopics() のデフォルト実装によって使用されます。
      パラメーター:
      createOrModifyTopic - 述語。
      導入:
      2.9.10
      関連事項:
    • getCreateOrModifyTopic

      protected PredicateSE<org.apache.kafka.clients.admin.NewTopic> getCreateOrModifyTopic()
      NewTopic を作成または変更するかどうかを決定するために使用される述語を返します。
      戻り値:
      述語。
      導入:
      2.9.10
      関連事項:
    • setClusterId

      public void setClusterId(StringSE clusterId)
      クラスター ID を設定します。これを使用すると、ユーザーが管理者権限を持っていない場合に、ブローカーからクラスター ID を取得しようとするのを防ぐことができます。
      パラメーター:
      clusterId - 設定する clusterId
      導入:
      3.1
    • getClusterId

      public @Nullable StringSE getClusterId()
      clusterId プロパティを取得します。
      戻り値:
      クラスター ID。
      導入:
      3.1.8
    • setClusterIdRetryInterval

      public void setClusterIdRetryInterval(DurationSE clusterIdRetryInterval)
      Set the minimum interval between attempts to fetch the cluster id from the broker after a failed attempt. Defaults to 5 minutes . Set to Duration.ZEROSE to attempt a fetch on every clusterId() call.
      パラメーター:
      clusterIdRetryInterval - 間隔。
      導入:
      4.1.1
      関連事項:
    • getConfigurationProperties

      public MapSE<StringSE,ObjectSE> getConfigurationProperties()
      インターフェースからコピーされた説明: KafkaAdminOperations
      この管理者の構成の変更不可能なコピーを取得します。
      次で指定:
      インターフェース KafkaAdminOperations 内の getConfigurationProperties 
      戻り値:
      構成マップ。
    • afterSingletonsInstantiated

      public void afterSingletonsInstantiated()
      次で指定:
      インターフェース SmartInitializingSingleton 内の afterSingletonsInstantiated 
    • initialize

      public final boolean initialize()
      このメソッドを呼び出して、トピックをチェック / 追加します。これは、アプリケーションコンテキストが初期化されたときにブローカーが使用できず、fatalIfBrokerNotAvailable が false であるか、autoCreate が false に設定されている場合に必要になることがあります。
      戻り値:
      成功した場合は true。
      関連事項:
    • newTopics

      protected CollectionSE<org.apache.kafka.clients.admin.NewTopic> newTopics()
      作成または変更する NewTopic のコレクションを返します。デフォルトの実装では、アプリケーションコンテキスト内のすべての NewTopic Bean を取得し、それぞれに setCreateOrModifyTopic(Predicate) 述語を適用します。同じトピック名の NewTopic がある場合は、TopicForRetryable Bean も削除されます。これは述語を呼び出す前に実行されます。
      戻り値:
      NewTopic のコレクション。
      導入:
      2.9.10
      関連事項:
    • clusterId

      public @Nullable StringSE clusterId()
      Return the cluster id, fetching it from the broker on the first call. A failed fetch is remembered and not attempted again until the clusterIdRetryInterval has elapsed; each attempt creates an Admin instance and blocks for up to operationTimeout, 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

      public void deleteTopics(StringSE... topicNames)
      Kafka クラスターからトピックを削除します。
      次で指定:
      インターフェース KafkaAdminOperations 内の deleteTopics 
      パラメーター:
      topicNames - 削除するトピック名。
      例外:
      KafkaException - 操作が失敗した場合。
      導入:
      4.0
    • createAdmin

      protected org.apache.kafka.clients.admin.Admin createAdmin()
      AdminClient クラスを使用して新しい Admin クライアントインスタンスを作成します。
      戻り値:
      新しい Admin クライアントインスタンス。
      導入:
      3.3.0
      関連事項:
      • AdminClient.create(Map)
    • getAdminConfig

      protected MapSE<StringSE,ObjectSE> getAdminConfig()