Skip to main content

AdminClient

Struct AdminClient 

Source
pub struct AdminClient { /* private fields */ }
Expand description

Kafka admin client for cluster administration.

Implementations§

Source§

impl AdminClient

Source

pub async fn describe_acls( &self, filter: AclFilter, ) -> Result<DescribeAclsResult>

Describe ACLs matching a filter.

§Example
ⓘ
// Describe all ACLs for a specific topic
let filter = AclFilter::for_resource(AclResourceType::Topic, "my-topic");
let result = admin.describe_acls(filter).await?;
Source

pub async fn create_acls( &self, acls: Vec<AclBinding>, ) -> Result<CreateAclsResult>

Create ACLs.

Returns Ok(result) when the RPC succeeds. An Ok return does not mean every ACL was created — inspect each element of CreateAclsResult::results for per-ACL failures.

§Arguments
  • acls - List of ACL bindings to create
§Example
ⓘ
let acl = AclBinding::allow_read_topic("my-topic", "User:alice");
admin.create_acls(vec![acl]).await?;
Source

pub async fn delete_acls( &self, filters: Vec<AclBindingFilter>, ) -> Result<DeleteAclsResult>

Delete ACLs matching the specified filters.

Returns Ok(result) when the RPC succeeds. An Ok return does not mean every filter matched or every ACL was deleted — inspect each element of DeleteAclsResult for per-filter failures.

§Arguments
  • filters - List of ACL binding filters to match for deletion
§Example
ⓘ
// Delete all ACLs for a specific topic
let filter = AclBindingFilter {
    resource_type: AclResourceType::Topic,
    resource_name: Some("my-topic".to_string()),
    pattern_type: AclPatternType::Literal,
    principal: None,
    host: None,
    operation: AclOperation::Any,
    permission_type: AclPermissionType::Any,
};
admin.delete_acls(vec![filter]).await?;
Source§

impl AdminClient

Source

pub async fn describe_configs_per_resource( &self, request: DescribeConfigsRequest, ) -> Result<Vec<DescribeConfigsResourceResult>>

Describe configuration for one or more resources (topics, brokers, etc.), preserving per-resource errors and attribution.

Uses DescribeConfigs (API Key 32). Build a DescribeConfigsRequest via its convenience constructors (for_topic, for_broker) or manually populate the resources field for multi-resource queries.

Each returned DescribeConfigsResourceResult carries the resource it describes and that resource’s error, if any. Use this rather than describe_configs whenever more than one resource is requested, or whenever an authorization failure must be distinguishable from “this resource has no configs”.

Source

pub async fn describe_configs( &self, request: DescribeConfigsRequest, ) -> Result<Vec<ConfigEntry>>

Describe configuration for one or more resources, flattening every resource’s entries into a single list.

§Errors

Unlike the previous behaviour, a per-resource error is not silently swallowed: if any requested resource failed (for example with TOPIC_AUTHORIZATION_FAILED), this returns that error rather than an empty or partial list that is indistinguishable from “no configs”.

Use describe_configs_per_resource to inspect partial results across a multi-resource request.

Source

pub async fn alter_topic_config( &self, topic: &str, configs: HashMap<String, String>, ) -> Result<AlterConfigResult>

Alter configuration for a topic.

Uses IncrementalAlterConfigs (API Key 44) to set individual config keys without replacing the entire config. Each key-value pair is applied as a SET operation.

IncrementalAlterConfigs is a controller-only API: the request is routed to the current controller and re-issued against a freshly resolved controller on NOT_CONTROLLER.

Source

pub async fn list_topics(&self) -> Result<Vec<String>>

List all topics.

Forces a metadata refresh first, so this is safe to call immediately after create_topics: the refresh waits out the retry.backoff.ms rate limiter rather than returning a stale cache as though it were current.

Source

pub async fn topic_config( &self, topic: &str, keys: &[&str], ) -> Result<HashMap<String, ConfigValue>>

Fetch specific config keys for a topic.

A convenience wrapper around describe_configs for the common case of reading a small set of well-known topic-level keys. Pass an empty keys slice to fetch all config entries for the topic.

Returns a map of config key → ConfigValue.

§Examples
ⓘ
let cfg = admin.topic_config("my-topic", &["retention.ms", "cleanup.policy"]).await?;
if let Some(krafka::admin::ConfigValue::Value(v)) = cfg.get("retention.ms") {
    println!("retention.ms = {v}");
}
Source

pub async fn describe_topics<S: AsRef<str>>( &self, topics: &[S], ) -> Result<HashMap<String, Arc<TopicInfo>>>

Describe topics by name, returning a map of topic name → [TopicInfo].

Topics not found in cluster metadata are absent from the returned map; callers can detect missing topics by comparing request and response key sets.

Source

pub async fn describe_topic( &self, topic: &str, ) -> Result<Option<Arc<TopicInfo>>>

Describe a single topic by name.

Returns None if the topic does not exist or is not visible in cluster metadata.

Source

pub async fn describe_cluster(&self) -> Result<DescribeClusterResult>

Describe the cluster using the DescribeCluster API (Key 60).

Returns cluster metadata including cluster ID, controller, brokers, and authorized operations.

Source

pub async fn partition_count(&self, topic: &str) -> Result<Option<usize>>

Get partition count for a topic.

Source

pub fn client_id(&self) -> &str

Get the client ID.

Source

pub fn request_timeout(&self) -> Duration

Get the request timeout.

Source§

impl AdminClient

Source

pub async fn describe_features(&self) -> Result<DescribeFeaturesResult>

Describe broker-supported and cluster-finalized features (KIP-584).

Sends an ApiVersions request (v3+) to any broker and extracts the feature information from the tagged fields. The response includes:

  • Features supported by the responding broker (per-broker)
  • Cluster-wide finalized features and their epoch (cluster-wide)
§Example
ⓘ
let features = admin.describe_features().await?;
for f in &features.supported_features {
    println!("{}: v{}–v{}", f.name, f.min_version, f.max_version);
}
for f in &features.finalized_features {
    println!("{}: v{}–v{} (finalized)", f.name, f.min_version_level, f.max_version_level);
}
Source

pub async fn update_features( &self, feature_updates: Vec<FeatureUpdateKey>, validate_only: bool, ) -> Result<UpdateFeaturesResult>

Update cluster-wide finalized feature version levels (KIP-584).

This is a destructive operation — downgrades and deletions can be data-lossy. Only the controller serves this request, so the client resolves the controller from cluster metadata and sends it there directly, re-resolving and retrying on NOT_CONTROLLER.

Requires ALTER permission on the cluster.

§Per-feature results are version-dependent

UpdateFeatures v2 (Kafka 4.0+) removed the per-feature Results array from the response: the controller now answers with a single top-level outcome. Against a v2 broker, UpdateFeaturesResult::results is therefore empty, and a successful return means every requested update was applied. Against a v0/v1 broker the per-feature breakdown is still populated. Treat Ok(_) — not a non-empty results — as the success signal, and use results only as extra diagnostics when it happens to be present.

§Example
ⓘ
use krafka::protocol::FeatureUpdateKey;

// `Ok` means the update was applied. `results` is a v0/v1-only detail.
let results = admin.update_features(
    vec![FeatureUpdateKey::upgrade("metadata.version", 17)],
    false, // validate_only
).await?;

for result in &results.results {
    if let Some(e) = &result.error {
        eprintln!("Failed to update {}: {e}", result.feature);
    }
}
Source§

impl AdminClient

Source

pub async fn delete_consumer_group_offsets( &self, group_id: &str, topic_partitions: &[(&str, &[i32])], ) -> Result<OffsetDeleteResult>

Delete committed offsets for a consumer group.

This is a destructive operation — deleted offsets cannot be recovered. The consumer group must be in the Empty state.

The request is sent to the group coordinator.

§Example
ⓘ
let results = admin.delete_offsets(
    "my-group",
    &[("my-topic", &[0, 1, 2])],
).await?;
Source

pub async fn describe_consumer_group_offsets( &self, group_id: &str, topic_partitions: Option<&[(&str, &[i32])]>, visibility: OffsetVisibility, ) -> Result<Vec<GroupOffsetEntry>>

Fetch committed offsets for a consumer group.

Pass topic_partitions to fetch offsets for specific partitions, or None to fetch all committed offsets for the group.

The request is sent to the group coordinator.

visibility decides what to do about an offset a transaction has staged but not yet committed — see OffsetVisibility. It is an explicit argument rather than a default because the two answers differ on exactly the pipelines where the difference matters, and a lag dashboard silently reading the unstable value reports progress that an abort can take back.

§Example
ⓘ
let offsets = admin
    .describe_consumer_group_offsets("my-group", None, OffsetVisibility::IncludeUnstable)
    .await?;
for entry in &offsets {
    println!("{}/{}: {}", entry.topic, entry.partition, entry.committed_offset);
}
§Errors

With OffsetVisibility::StableOnly, a partition whose offset is staged inside an unresolved transaction is reported by the broker as UNSTABLE_OFFSET_COMMIT; it is surfaced rather than omitted, since an omitted partition is indistinguishable from one the group never committed.

Source

pub async fn alter_consumer_group_offsets( &self, group_id: &str, topic_offsets: &[(&str, &[(i32, i64)])], ) -> Result<Vec<AlterGroupOffsetResult>>

Alter committed offsets for a consumer group.

Sets each specified partition’s committed offset. The consumer group must be in the Empty state (no active members).

The request is sent to the group coordinator.

§Example
ⓘ
admin
    .alter_consumer_group_offsets(
        "my-group",
        &[("my-topic", &[(0, 100), (1, 200)])],
    )
    .await?;
Source§

impl AdminClient

Source

pub async fn describe_consumer_groups( &self, group_ids: Vec<String>, ) -> Result<Vec<ConsumerGroupDescription>>

Describe consumer groups.

Automatically detects whether each group uses the classic protocol or the new consumer protocol (KIP-848) and dispatches to the appropriate API:

  • Classic groups → DescribeGroups (Key 15)
  • Consumer groups → ConsumerGroupDescribe (Key 69)

The returned ConsumerGroupDescription is a unified type. Fields specific to one protocol variant are Option-wrapped.

§Example
ⓘ
let groups = admin
    .describe_consumer_groups(vec!["my-group".to_string()])
    .await?;
for group in &groups {
    println!("{}: type={}, state={}, members={}",
        group.group_id, group.group_type, group.state, group.members.len());
}
Source

pub async fn list_consumer_groups( &self, filter: &GroupListing, ) -> Result<Vec<ConsumerGroupListing>>

List all consumer groups on the cluster.

Returns a list of all consumer groups with their protocol types.

§Example
ⓘ
let groups = admin.list_consumer_groups(&GroupListing::all()).await?;
for group in &groups {
    println!("{} ({})", group.group_id, group.protocol_type);
}

filter is applied by the broker. On a cluster with thousands of groups that is the difference between transferring the whole registry and transferring the handful you asked about — see GroupListing.

Source

pub async fn delete_records( &self, offsets: HashMap<(String, i32), i64>, timeout: Duration, ) -> Result<Vec<DeleteRecordResult>>

Delete records from topic partitions before the specified offsets.

Records with offsets less than the specified offset for each partition will be marked for deletion. This adjusts the log start offset.

§Arguments
  • offsets - Map of (topic, partition) to the offset before which to delete
  • timeout - Operation timeout
§Example
ⓘ
use std::collections::HashMap;
let mut offsets = HashMap::new();
offsets.insert(("my-topic".to_string(), 0), 100i64);
let results = admin.delete_records(offsets, Duration::from_secs(30)).await?;
Source

pub async fn offset_for_leader_epoch( &self, partitions: Vec<(String, i32, i32)>, ) -> Result<Vec<LeaderEpochResult>>

Get the end offset for each partition at the given leader epoch.

This is used to detect log truncation after a leader change. For each topic-partition, the broker returns the end offset for the requested leader epoch. If the epoch is no longer valid, the broker returns the epoch and offset where the log was truncated.

§Arguments
  • partitions - List of (topic, partition, leader_epoch) tuples
§Example
ⓘ
let results = admin.offset_for_leader_epoch(
    vec![("my-topic".to_string(), 0, 5)]
).await?;
for r in &results {
    println!("{}:{} epoch={} end_offset={}", r.topic, r.partition, r.leader_epoch, r.end_offset);
}
Source§

impl AdminClient

Source

pub async fn list_offsets( &self, topic_partitions: &[(&str, &[i32])], spec: OffsetSpec, ) -> Result<Vec<ListOffsetResult>>

List offsets for one or more topic-partitions.

Each request is routed to the partition’s current leader. Metadata is refreshed once on NotLeaderForPartition errors before retrying.

§Arguments
  • topic_partitions — slice of (topic_name, partition_ids) pairs.
  • spec — which offset to fetch (Earliest, Latest, or Timestamp).
§Example
ⓘ
use krafka::admin::{AdminClient, OffsetSpec};

let results = admin
    .list_offsets(&[("my-topic", &[0, 1, 2])], OffsetSpec::Latest)
    .await?;
for r in &results {
    println!("{}/{}: offset={}", r.topic, r.partition, r.offset);
}
Source

pub async fn consumer_group_lag( &self, group_id: &str, topic_partitions: Option<&[(&str, &[i32])]>, ) -> Result<Vec<ConsumerGroupLag>>

Compute consumer group lag for the specified topics.

Lag is defined as end_offset − committed_offset for each topic-partition. Partitions with no committed offset have committed_offset = None and lag = None.

Partitions whose end offset could not be fetched report end_offset = None, lag = None, and the reason in ConsumerGroupLag::end_offset_error — never lag = 0, which would make a stalled consumer look healthy.

This method issues two parallel-ish requests:

  1. describe_consumer_group_offsets for the committed positions.
  2. list_offsets with OffsetSpec::Latest for the end offsets.

The consumer group does not need to be stopped.

§Arguments
  • group_id — consumer group ID.
  • topic_partitions — which partitions to measure; pass None to measure all partitions that the group has committed offsets for.
§Example
ⓘ
let lag = admin
    .consumer_group_lag("my-group", Some(&[("my-topic", &[0, 1, 2])]))
    .await?;
for entry in &lag {
    println!(
        "{}/{}: lag={:?}",
        entry.topic, entry.partition, entry.lag
    );
}
Source§

impl AdminClient

Source

pub async fn describe_log_dirs( &self, topics: Option<Vec<DescribableLogDirTopic>>, ) -> Result<Vec<LogDirInfo>>

Describe log directories on all known brokers.

Each broker maintains one or more log directories; this method queries every broker and returns per-directory information including sizes, partition assignments, and (v4+) volume capacity.

Pass None for topics to describe all partitions on every broker, or pass a list of [DescribableLogDirTopic] to filter.

§Example
ⓘ
// Describe all log dirs on every broker
let dirs = admin.describe_log_dirs(None).await?;
for dir in &dirs {
    println!("broker {} dir {}: {:?}", dir.broker_id, dir.log_dir, dir.error);
}

// Describe specific topic partitions
use krafka::protocol::DescribableLogDirTopic;
let filter = vec![DescribableLogDirTopic {
    topic: "my-topic".into(),
    partitions: vec![0, 1, 2],
}];
let dirs = admin.describe_log_dirs(Some(filter)).await?;
Source

pub async fn elect_leaders( &self, election_type: ElectionType, topic_partitions: Option<Vec<ElectLeadersTopicPartitions>>, timeout: Duration, ) -> Result<Vec<ElectLeadersResult>>

Trigger leader election for the specified partitions.

When topic_partitions is None, leaders for all partitions are elected. The election_type controls whether to perform a preferred or unclean leader election (requires broker v1+; v0 always does preferred election).

Returns per-partition results — individual partitions may fail even when the RPC succeeds.

§Example
ⓘ
use krafka::protocol::ElectionType;
use std::time::Duration;

// Preferred election for all partitions
let results = admin
    .elect_leaders(ElectionType::Preferred, None, Duration::from_secs(60))
    .await?;
Source

pub async fn alter_partition_reassignments( &self, topics: Vec<ReassignableTopic>, timeout: Duration, ) -> Result<AlterReassignmentsResult>

Alter partition reassignments.

Initiates or cancels partition reassignments. To cancel a pending reassignment, set replicas to None for that partition.

This is a destructive operation — reassigning partitions moves data between brokers and can significantly impact cluster load.

Returns per-partition results — individual partitions may fail even when the RPC succeeds.

§Example
ⓘ
use krafka::protocol::{ReassignableTopic, ReassignablePartition};
use std::time::Duration;

let results = admin.alter_partition_reassignments(
    vec![ReassignableTopic {
        name: "my-topic".into(),
        partitions: vec![ReassignablePartition {
            partition_index: 0,
            replicas: Some(vec![1, 2, 3]),
        }],
    }],
    Duration::from_secs(60),
).await?;
Source

pub async fn alter_partition_reassignments_opts( &self, topics: Vec<ReassignableTopic>, timeout: Duration, allow_replication_factor_change: bool, ) -> Result<AlterReassignmentsResult>

alter_partition_reassignments with control over whether the move may change a partition’s replication factor (AlterPartitionReassignments v1, Kafka 4.1).

Passing false asks the broker to reject any target replica set whose size differs from the partition’s current replica count, answering with INVALID_REPLICATION_FACTOR instead of applying the change. That is the safer default for generated reassignment plans, where an off-by-one would otherwise silently reduce durability.

§Version requirement

The field only exists from v1. Against a v0 broker allow_replication_factor_change = false cannot be honoured, so this returns ProtocolErrorKind::UnknownApiVersion rather than quietly sending a request that permits the change. Passing true — the broker default and the v0 behaviour — works against every version.

Source

pub async fn list_partition_reassignments( &self, topics: Option<Vec<ListPartitionReassignmentsTopic>>, timeout: Duration, ) -> Result<Vec<PartitionReassignmentInfo>>

List ongoing partition reassignments.

When topics is None, all ongoing reassignments are listed. Otherwise, only the specified topic-partitions are checked.

§Example
ⓘ
// List all ongoing reassignments
let reassignments = admin
    .list_partition_reassignments(None, Duration::from_secs(60))
    .await?;
for topic in &reassignments {
    for p in &topic.partitions {
        println!("{} p{}: adding {:?}, removing {:?}",
            topic.name, p.partition_index, p.adding_replicas, p.removing_replicas);
    }
}
Source

pub async fn alter_replica_log_dirs( &self, dirs: Vec<AlterReplicaLogDir>, ) -> Result<Vec<AlterReplicaLogDirsResult>>

Move partition replicas to a different log directory on the broker.

This is a destructive operation — moving replicas between log directories triggers data copying and can impact broker I/O.

This is a per-broker operation. Each broker is sent only the topic-partitions it actually hosts a replica for, resolved from cached cluster metadata. Brokers that host none of the requested replicas are not contacted at all.

This matters: broadcasting the full request to every broker asks each one to move replicas it does not own, producing a storm of REPLICA_NOT_AVAILABLE / KAFKA_STORAGE_ERROR results that are indistinguishable from genuine failures on the brokers that do own them.

If metadata knows no replicas for a requested partition, the request is sent to the partition’s leader when known, and reported as an error otherwise — it is never broadcast.

Returns per-partition results — individual partitions may fail even when the RPC succeeds.

§Example
ⓘ
use krafka::protocol::{AlterReplicaLogDir, AlterReplicaLogDirTopic};

let results = admin.alter_replica_log_dirs(vec![
    AlterReplicaLogDir {
        path: "/data/kafka-logs-2".into(),
        topics: vec![AlterReplicaLogDirTopic {
            name: "my-topic".into(),
            partitions: vec![0, 1],
        }],
    },
]).await?;
Source§

impl AdminClient

Source

pub async fn describe_client_quotas( &self, components: &[(&str, i8, Option<&str>)], strict: bool, ) -> Result<DescribeClientQuotasResult>

Describe client quotas matching the given filter.

§Arguments
  • components - Filter components. Each component specifies an entity type and match criteria. The broker returns entities matching all components.
  • strict - If true, exclude entities with unspecified entity types (i.e., only return entities that exactly match all given component types).
§Filter Match Types

Each component has a match_type:

  • 0 (exact): match the entity with the given name
  • 1 (default): match the default entity for this type
  • 2 (any specified): match any entity with a name (non-default)
§Example
ⓘ
// Describe all quotas for user "alice"
let results = admin.describe_client_quotas(
    &[("user", 0, Some("alice"))],
    false,
).await?;
Source

pub async fn alter_client_quotas( &self, entries: &[QuotaAlteration<'_>], validate_only: bool, ) -> Result<Vec<AlterClientQuotaResult>>

Alter client quotas.

Each entry specifies an entity (user, client-id, ip) and a set of quota operations (set or remove). Results are returned per-entity.

§Arguments
  • entries - Quota alterations. Each entry has an entity and operations.
  • validate_only - If true, validate the request without applying changes.

AlterClientQuotas is a controller-only API: the request is routed to the current controller and re-issued against a freshly resolved controller on NOT_CONTROLLER.

§Example
ⓘ
use krafka::admin::QuotaAlteration;

// Set producer byte rate quota for user "alice"
let results = admin.alter_client_quotas(
    &[QuotaAlteration {
        entity: vec![("user", Some("alice"))],
        ops: vec![("producer_byte_rate", Some(1_048_576.0))],
    }],
    false,
).await?;
Source§

impl AdminClient

Source

pub async fn describe_user_scram_credentials( &self, users: Option<Vec<String>>, ) -> Result<DescribeUserScramCredentialsResult>

Describe SCRAM credentials for the specified users.

When users is None, all SCRAM credentials are described.

§Example
ⓘ
// Describe all SCRAM credentials
let results = admin.describe_user_scram_credentials(None).await?;
for user in &results {
    println!("{}: {:?}", user.name, user.credential_infos);
}
Source

pub async fn alter_user_scram_credentials( &self, deletions: Vec<ScramCredentialDeletion>, upsertions: Vec<ScramCredentialUpsertion>, ) -> Result<Vec<AlterScramCredentialResult>>

Alter (upsert or delete) SCRAM credentials for users.

This is a destructive operation — deleting a SCRAM credential removes the user’s ability to authenticate with that mechanism.

AlterUserScramCredentials is a controller-only API: the request is routed to the current controller and re-issued against a freshly resolved controller on NOT_CONTROLLER.

§Example
ⓘ
use krafka::protocol::{ScramCredentialDeletion, ScramCredentialUpsertion};
use krafka::auth::ScramMechanism;
use zeroize::Zeroizing;

let results = admin.alter_user_scram_credentials(
    vec![ScramCredentialDeletion {
        name: "alice".into(),
        mechanism: ScramMechanism::Sha512,
    }],
    vec![ScramCredentialUpsertion {
        name: "bob".into(),
        mechanism: ScramMechanism::Sha256,
        iterations: 8192,
        salt: Zeroizing::new(vec![1, 2, 3]),
        salted_password: Zeroizing::new(vec![4, 5, 6]),
    }],
).await?;
Source§

impl AdminClient

Source

pub async fn describe_share_group_offsets( &self, group_id: &str, topics: Option<&[(&str, &[i32])]>, ) -> Result<DescribeShareGroupOffsetsResult>

Read a share group’s share-partition start offsets (KIP-932).

Pass topics = None to describe every topic-partition the group holds state for. Passing Some(&[]) describes nothing — the wire protocol distinguishes a null topics array from an empty one, and so does this method.

lag is Some(_) only when the coordinator supports DescribeShareGroupOffsets v1 (KIP-1226, Kafka 4.3+); against an older broker it is None rather than a misleading zero.

The request goes to the group coordinator, which is where share-group state lives.

§Example
ⓘ
// Everything the group knows about.
let all = admin.describe_share_group_offsets("orders-share", None).await?;
for p in &all.partitions {
    println!("{}-{} start={} lag={:?}", p.topic, p.partition, p.start_offset, p.lag);
}

// Just two partitions of one topic.
let some = admin
    .describe_share_group_offsets("orders-share", Some(&[("orders", &[0, 1][..])]))
    .await?;
§Errors

Returns an error if the client is closed, a topic name is invalid, the coordinator cannot be found, or the broker supports no compatible DescribeShareGroupOffsets version.

Source

pub async fn alter_share_group_offsets( &self, group_id: &str, topic_offsets: &[(&str, &[(i32, i64)])], ) -> Result<Vec<ShareGroupOffsetAlteration>>

Reset a share group’s share-partition start offsets (KIP-932).

This is a destructive operation. Moving the start offset backwards re-delivers records the group already processed; moving it forwards skips records permanently.

The group must be empty. A group with a live member is answered with NON_EMPTY_GROUP, for the same reason alter_consumer_group_offsets requires it: rewriting the start offset under an active member would hand it records it has already acquired.

§Example
ⓘ
// Rewind two partitions to the beginning of the log.
let results = admin
    .alter_share_group_offsets("orders-share", &[("orders", &[(0, 0), (1, 0)][..])])
    .await?;
for r in &results {
    if let Some(e) = &r.error {
        eprintln!("{}-{} failed: {e}", r.topic, r.partition);
    }
}
§Errors

Returns an error if the client is closed, a topic name is invalid, the coordinator cannot be found, the broker supports no compatible AlterShareGroupOffsets version, or the request fails at the top level (including NON_EMPTY_GROUP).

Source

pub async fn delete_share_group_offsets( &self, group_id: &str, topics: &[&str], ) -> Result<Vec<ShareGroupOffsetDeletion>>

Delete a share group’s offset state for whole topics (KIP-932).

This is a destructive operation — the deleted state cannot be recovered, and the group restarts those topics from its configured reset policy. The group must be empty.

Use it after retiring a topic, so the coordinator stops carrying state for partitions that no longer exist.

§Errors

Returns an error if the client is closed, a topic name is invalid, the coordinator cannot be found, the broker supports no compatible DeleteShareGroupOffsets version, or the request fails at the top level (including NON_EMPTY_GROUP).

Source§

impl AdminClient

Source

pub async fn describe_streams_groups( &self, group_ids: &[&str], ) -> Result<Vec<DescribedStreamsGroup>>

Describe Streams groups (KIP-1071, Kafka 4.1+).

Returns each group’s topology, members, and per-member task assignments and changelog offsets. Like every group API, the request goes to each group’s coordinator, so groups spread across brokers are batched per coordinator and issued once per broker.

§What this is for

krafka has no Streams runtime and cannot join a Streams group — that is StreamsGroupHeartbeat (key 88), whose request carries the application topology. This is the observational half: it is what an operator needs to answer “is this Streams application healthy?” without running one.

Two fields are the ones worth alerting on:

  • A member whose topology_epoch is below the group’s StreamsTopology::epoch is still running an older topology.
  • A member whose assignment differs from its target_assignment has not finished rebalancing. Persistently so usually means restoration is not keeping up — compare task_offsets against task_end_offsets for the lag that explains it.
§Errors

Per-group failures are reported in each DescribedStreamsGroup::error_code rather than failing the whole call, so one unknown group does not hide the others. A broker that does not support the API at all — anything before Kafka 4.1 — fails the call with ProtocolErrorKind::UnknownApiVersion.

§Example
ⓘ
let groups = admin.describe_streams_groups(&["my-streams-app"]).await?;
for group in &groups {
    if !group.error_code.is_ok() {
        eprintln!("{}: {:?}", group.group_id, group.error_code);
        continue;
    }
    let topology_epoch = group.topology.as_ref().map_or(-1, |t| t.epoch);
    for member in &group.members {
        let lagging = member.topology_epoch < topology_epoch;
        let rebalancing = member.assignment != member.target_assignment;
        println!(
            "{} active={} lagging_topology={lagging} rebalancing={rebalancing}",
            member.member_id,
            member.assignment.active_tasks.len(),
        );
    }
}
Source§

impl AdminClient

Source

pub async fn create_delegation_token( &self, owner: Option<(&str, &str)>, renewers: &[(&str, &str)], max_lifetime: Option<Duration>, ) -> Result<CreateDelegationTokenResult>

Create a delegation token.

Delegation tokens allow a principal to delegate authentication to another principal without sharing credentials (KIP-48). The token HMAC can be used for SASL/SCRAM authentication.

CreateDelegationToken is a controller-only API: the request is routed to the current controller and re-issued against a freshly resolved controller on NOT_CONTROLLER.

§Arguments
  • owner - Principal the token is issued for, as a (type, name) pair, or None for the authenticated caller. Requesting a token on behalf of another principal is KIP-373 and needs CreateDelegationToken v3+ plus CreateTokens authorisation on that principal — it is how a superuser provisions a token for a service account that never authenticates interactively. The requester is recorded alongside the owner and comes back on DelegationToken::token_requester_principal_name, which is what an audit trail needs.
  • renewers - Principals authorized to renew the token (type, name pairs). Pass an empty slice to allow only the token owner to renew.
  • max_lifetime - Maximum token lifetime. Use None for the server default (typically 7 days).
Source

pub async fn renew_delegation_token( &self, hmac: &[u8], renew_period: Duration, ) -> Result<RenewDelegationTokenResult>

Renew a delegation token, extending its expiry time.

§Arguments
  • hmac - HMAC of the token to renew (from DelegationToken::hmac).
  • renew_period - How long to extend the token’s lifetime.
Source

pub async fn expire_delegation_token( &self, hmac: &[u8], expiry_period: Option<Duration>, ) -> Result<ExpireDelegationTokenResult>

Expire a delegation token, revoking it before its natural expiry.

§Arguments
  • hmac - HMAC of the token to expire (from DelegationToken::hmac).
  • expiry_period - How long until the token expires. Pass None to expire the token immediately (sends -1 to the broker).
Source

pub async fn describe_delegation_token( &self, owners: Option<&[(&str, &str)]>, ) -> Result<Vec<DelegationToken>>

Describe delegation tokens visible to the caller.

§Arguments
  • owners - Filter by token owners (type, name pairs). Pass None to return all tokens visible to the caller.
Source§

impl AdminClient

Source

pub async fn create_topics( &self, topics: Vec<NewTopic>, timeout: Duration, validate_only: bool, ) -> Result<Vec<CreateTopicResult>>

Create topics.

CreateTopics is a controller-only API. The request is routed to the current controller and re-issued against a freshly resolved controller if the broker answers NOT_CONTROLLER — see get_controller_connection. A NOT_CONTROLLER that survives every retry is returned as an Err, not as an Ok carrying a per-topic error string.

Returns Ok(results) when the RPC succeeds. An Ok return does not mean every topic was created — inspect each CreateTopicResult::error for per-topic failures. Every requested topic is guaranteed to appear in the result: a topic the broker omits from its response is reported with an explicit error rather than silently vanishing.

§Parameters
  • topics — Descriptions of the topics to create.
  • timeout — How long the broker should wait for the creation to complete.
  • validate_only — When true, the broker validates the request but does not create any topics. Useful for pre-flight checks. Requires CreateTopics v2+ (Kafka 0.11+); all modern brokers support this.
Source

pub async fn delete_topics( &self, topics: Vec<String>, timeout: Duration, ) -> Result<Vec<DeleteTopicResult>>

Delete topics.

DeleteTopics is a controller-only API; see create_topics for how controller routing and NOT_CONTROLLER retries work.

Returns Ok(results) when the RPC succeeds. An Ok return does not mean every topic was deleted — inspect each DeleteTopicResult::error for per-topic failures. Topics omitted by the broker are reported with an explicit error.

Source

pub async fn create_partitions( &self, topic: impl Into<String>, new_total_count: i32, timeout: Duration, validate_only: bool, ) -> Result<CreatePartitionsResult>

Increase the number of partitions for a topic.

Note: Partition count can only be increased, never decreased.

CreatePartitions is a controller-only API; see create_topics for how controller routing and NOT_CONTROLLER retries work.

§Parameters
  • topic — Topic to expand.
  • new_total_count — The total partition count after the increase.
  • timeout — How long the broker should wait for the change to complete.
  • validate_only — When true, the broker validates the request but does not create any partitions.
Source§

impl AdminClient

Source

pub async fn describe_producers( &self, topic_partitions: &[(&str, &[i32])], ) -> Result<Vec<DescribeProducersTopicResult>>

Describe active producers on the given topic-partitions.

Routes each topic-partition to its leader broker via cached metadata for optimal performance. Falls back to any broker if the leader is unknown.

Returns per-partition producer state useful for debugging transactional and idempotent producers.

§Example
ⓘ
let results = admin
    .describe_producers(&[("my-topic", &[0, 1])])
    .await?;
Source

pub async fn describe_transactions( &self, transactional_ids: &[&str], ) -> Result<Vec<TransactionDescription>>

Describe the state of the given transactions.

Routes each transactional ID to its transaction coordinator via FindCoordinator, groups by coordinator, and batches requests.

§Example
ⓘ
let results = admin
    .describe_transactions(&["txn-1", "txn-2"])
    .await?;
Source

pub async fn list_transactions( &self, state_filters: &[&str], producer_id_filters: &[i64], duration_filter: i64, transactional_id_pattern: Option<&str>, ) -> Result<ListTransactionsResult>

List transactions matching the given filters.

Queries all brokers and merges results, because each broker only knows about transactions it coordinates.

Pass empty slices for state_filters and producer_id_filters, -1 for duration_filter and None for transactional_id_pattern to list all transactions.

transactional_id_pattern requires ListTransactions v2 (KIP-1152, Kafka 4.1+). Against an older broker the negotiated version is lower and the pattern cannot be sent, so this method rejects the call rather than silently returning unfiltered results.

§Example
ⓘ
// List all ongoing transactions
let txns = admin
    .list_transactions(&["Ongoing"], &[], -1, None)
    .await?;

// Only transactions whose transactional ID starts with "orders-"
let txns = admin
    .list_transactions(&[], &[], -1, Some("orders-.*"))
    .await?;
Source

pub async fn list_config_resources( &self, resource_types: &[ConfigResourceType], ) -> Result<Vec<ListedConfigResource>>

List config resources known to the cluster (KIP-1142).

Kafka 4.1 renamed API key 74 from ListClientMetricsResources to ListConfigResources and generalised it: v0 could only enumerate client-metrics subscriptions (KIP-714), while v1 enumerates any config resource type — topics, brokers, broker loggers, groups and client metrics.

Pass an empty resource_types slice to get whichever types the broker lists by default.

§Version behaviour

Against a pre-4.1 broker only v0 is available, which ignores resource_types and returns client-metrics subscriptions. Requesting a type other than ConfigResourceType::ClientMetrics there would silently return the wrong set, so this method rejects the call instead.

§Example
ⓘ
use krafka::protocol::ConfigResourceType;

// Everything the broker lists by default.
let all = admin.list_config_resources(&[]).await?;

// Just the topics.
let topics = admin
    .list_config_resources(&[ConfigResourceType::Topic])
    .await?;
Source

pub async fn list_client_metrics_resources(&self) -> Result<Vec<String>>

List client metrics subscription names (KIP-714).

Convenience wrapper over list_config_resources that asks for ConfigResourceType::ClientMetrics and returns only the names. Works against any broker: this is exactly what API key 74 did before KIP-1142 renamed it.

§Example
ⓘ
let names = admin.list_client_metrics_resources().await?;
for name in &names {
    println!("subscription: {name}");
}
Source

pub async fn write_txn_markers( &self, markers: &[WritableTxnMarker], ) -> Result<Vec<WriteTxnMarkersResult>>

Write transaction markers (COMMIT or ABORT) to the given topic-partitions.

This is an inter-broker API used to finalize transactions. The admin client exposes it primarily for aborting stuck transactions (abort_transaction).

§Routing

WriteTxnMarkers must reach the leader of each partition: only the leader can append the marker to the log. This method groups each marker’s topic-partitions by their current leader and sends one WriteTxnMarkers request per leader, then merges the per-partition results back together.

Sending the whole marker set to a single broker instead would leave a transaction that spans several brokers only partially finalised: the other leaders answer NOT_LEADER_OR_FOLLOWER, the transaction stays open, and read_committed consumers block on the last stable offset indefinitely.

§Errors

Returns Err if a partition’s leader is unknown, or if the request to a leader fails outright — a marker that was never written must not be reported as success. Per-partition broker errors are surfaced in WriteTxnMarkersPartitionResult::error.

§Example
ⓘ
use krafka::protocol::{WritableTxnMarker, WritableTxnMarkerTopic};

let results = admin
    .write_txn_markers(&[WritableTxnMarker {
        producer_id: 42,
        producer_epoch: 5,
        transaction_result: false, // ABORT
        topics: vec![WritableTxnMarkerTopic {
            name: "my-topic".into(),
            partition_indexes: vec![0, 1],
        }],
        coordinator_epoch: 10,
        transaction_version: 0,
    }])
    .await?;
Source

pub async fn abort_transaction( &self, transactional_id: &str, ) -> Result<Vec<WriteTxnMarkersResult>>

Abort a stuck transaction by writing an ABORT marker.

This is the admin-friendly wrapper around write_txn_markers that discovers the affected partitions via describe_transactions, determines the owning coordinator epoch via describe_producers, and writes an ABORT marker to each partition’s leader.

§The coordinator epoch is discovered, never assumed

The partition leader validates the marker’s coordinator_epoch against the epoch it has cached for that producer, and accepts any epoch when it has none cached. Passing a fabricated 0 therefore succeeds exactly on the partitions where validation cannot protect you — aborting a transaction that a live, newer coordinator still owns and discarding committed exactly-once data.

This method instead reads the real epoch from the active producer state on the target partitions, mirroring Java’s AbortTransactionSpec. abort_transaction_with_epoch takes the epoch explicitly when the caller already knows it.

§Errors

Returns an error when the transaction cannot be described, when no active producer state exposes a coordinator epoch (so no safe epoch can be determined), or when the discovered producer state disagrees across partitions.

§Example
ⓘ
admin.abort_transaction("my-transactional-id").await?;
Source

pub async fn abort_transaction_with_epoch( &self, transactional_id: &str, coordinator_epoch: i32, ) -> Result<Vec<WriteTxnMarkersResult>>

Abort a stuck transaction using a caller-supplied coordinator epoch.

Use this when the epoch is already known (for example from a previous describe_producers call); otherwise prefer abort_transaction, which discovers it.

§Warning

Supplying an epoch older than the one the partition leader has cached is rejected by the broker, but supplying an arbitrary epoch to a leader with no cached epoch is accepted. Never invent a value here.

Source

pub async fn describe_metadata_quorum( &self, topic_partitions: &[(&str, &[i32])], ) -> Result<DescribeQuorumResult>

Describe the KRaft quorum for the given topic-partitions.

In a KRaft-mode cluster this returns the current voters, observers, leader, leader epoch, and high watermark for each quorum partition.

The primary use case is inspecting __cluster_metadata partition 0.

§Example
ⓘ
let result = admin
    .describe_metadata_quorum(&[("__cluster_metadata", &[0])])
    .await?;
Source§

impl AdminClient

Source

pub fn builder() -> AdminClientBuilder

Create a new admin client builder.

Source

pub async fn get_controller_connection(&self) -> Result<Arc<BrokerConnection>>

Get a connection to the cluster controller.

Controller-only APIs must be routed here rather than to an arbitrary broker. A non-controller broker does forward such requests, but during a controller failover it answers NOT_CONTROLLER (41) instead — which the admin APIs surface only as a per-item error string, so a caller checking nothing but the Result concludes the operation succeeded.

§Errors

Returns crate::error::ErrorCode::UnknownControllerId if the cluster reports no controller even after a metadata refresh.

Source

pub async fn delete_consumer_groups( &self, group_ids: Vec<String>, ) -> Result<Vec<DeleteGroupResult>>

Delete consumer groups by ID.

Returns one DeleteGroupResult per group. Each result may contain an error if that particular group could not be deleted (e.g., it has active members).

Source

pub async fn describe_topic_partitions( &self, topics: Vec<String>, ) -> Result<DescribeTopicPartitionsResult>

Describe topic partitions using the DescribeTopicPartitions API (Key 75).

Returns detailed per-partition information including leader, replicas, ISR, eligible leader replicas (ELR), and offline replicas. Supports pagination for topics with many partitions.

§Example
ⓘ
let result = admin
    .describe_topic_partitions(vec!["my-topic".to_string()])
    .await?;
for topic in &result.topics {
    println!("{}: {} partitions", topic.name.as_deref().unwrap_or("?"), topic.partitions.len());
    for p in &topic.partitions {
        println!("  partition {}: leader={}, isr={:?}", p.partition_index, p.leader_id, p.isr_nodes);
    }
}
Source

pub fn pool(&self) -> &Arc<ConnectionPool> ⓘ

Get access to the connection pool.

Source

pub fn update_seed_brokers(&self, servers: Vec<String>) -> Result<()>

Replace the bootstrap server list at runtime (KIP-899).

The new addresses are used on the next metadata refresh that falls back to bootstrap servers. Does not close existing connections.

§Errors

Returns an error if servers is empty.

Source

pub async fn refresh_tls(&self) -> Result<()>

Re-read TLS certificate and key files from disk and atomically install the new material for all future connections (KIP-1288).

Existing TLS sessions are unaffected: they keep the connector they handshaked with and are replaced naturally as connections cycle. On error the previously loaded certificates stay active, so a call made mid-rotation against a half-written PEM is safe to retry.

No-op when TLS is not configured.

Use this for event-driven rotation (an inotify watch, a sidecar signal). For unattended rotation set TransportConfig::tls_reload_interval instead and krafka reloads on a timer.

§Errors

Returns an error if the certificate or key files cannot be read or parsed.

Source

pub async fn rebootstrap(&self)

Force a rebootstrap: close all connections, clear the metadata cache, and fall back to bootstrap servers (KIP-899).

Source

pub async fn close(&self)

Close the admin client.

Sets the closed flag so that subsequent operations fail fast.

§Connection teardown depends on pool ownership
  • Own pool (the client was built from bootstrap_servers): all broker connections are torn down. In-flight admin RPCs that have not yet received a response will fail with a network error, so callers should let long-running admin operations finish first.
  • Shared pool (built via AdminClientBuilder::with_client): connections are left untouched. The pool belongs to the KrafkaClient, and closing it here would kill every producer and consumer connection on that client and fail all in-flight Produce/Fetch requests. Close the KrafkaClient to release those connections.

Calling close() more than once is a no-op.

Source

pub fn owns_pool(&self) -> bool

Whether this admin client owns its connection pool.

false when the pool was borrowed from a KrafkaClient via AdminClientBuilder::with_client; in that case close does not tear down connections.

Source

pub fn is_closed(&self) -> bool

Check if the admin client is closed.

Source

pub fn connection_metrics(&self) -> Arc<ConnectionMetrics> ⓘ

Get the shared connection metrics handle used by this admin client’s broker pool.

Trait Implementations§

Source§

impl Drop for AdminClient

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<Unshared, Shared> IntoShared<Shared> for Unshared
where Shared: FromUnshared<Unshared>,

Source§

fn into_shared(self) -> Shared

Creates a shared type from an unshared type.
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more