---
title: "From Alert to Fix in 60 Seconds: AI-Powered Kafka Operations - Calinora Blog"
description: "How Pilot's MCP-powered AI agent traced slow order processing to a disk-hot broker, rebalanced the affected topic, and resolved consumer lag - all through natural language with human approval gates."
url: "https://blog.calinora.io/posts/ai-powered-kafka-operations-with-mcp/"
date: "2026-03-12"
---

March 12, 2026

# From Alert to Fix in 60 Seconds: AI-Powered Kafka Operations

It’s Wednesday afternoon. The payments team pings your on-call channel: “orders are processing slowly.” You open a terminal, SSH into a jumpbox, run `kafka-consumer-groups.sh --describe`, and stare at a wall of partition offsets trying to figure out which partitions are lagging and why.

Twenty minutes later, you’ve found the consumer group has lag. But is it a slow consumer? A hot broker? A topic partition imbalance? You need to cross-reference consumer lag with partition assignments, broker disk usage, and activity metrics. That’s another three CLI commands, a spreadsheet, and a lot of context-switching.

There’s a better way.

## The Scenario

Pilot ships with a built-in AI agent that connects directly to your Kafka cluster through the [Model Context Protocol (MCP)](https://modelcontextprotocol.io). Instead of memorizing CLI flags and cross-referencing outputs, you describe what you’re seeing in natural language. The agent has access to over 60 specialized tools - from consumer group inspection to broker health checks to partition rebalancing.

Here’s a real interaction. Press play to watch.

1. **Operator:** the payments team says orders are processing slowly. can you check their consumer group?
2. **Tool call** `describe_consumer_group`: group: "payments-processor"
3. **Tool call** `get_consumer_group_lag`: group: "payments-processor"
4. **Tool result** `describe_consumer_group`: group: payments-processor; state: Stable; members: 3; topics: payments-events
5. **Tool result** `get_consumer_group_lag`: totalLag: 847,231; partitions: 12/12 assigned; maxLag: P3: 312K, P7: 289K; avgLag: 70.6K
6. **Agent:** payments-processor is Stable with 3 members, but has 847K messages of total lag - concentrated on partitions 3 and 7 which account for 71% of all lag. That pattern suggests an infrastructure issue rather than a slow consumer. Let me dig into those partitions.
7. **Tool call** `get_partition_activity`: topic: "payments-events"
8. **Tool call** `get_cluster_health`: Checking broker health...
9. **Tool result** `get_partition_activity`: P3 leader: broker-4; P7 leader: broker-4; broker-4 leaders: 48% of topic; throughput: 1,240 msg/s
10. **Tool result** `get_cluster_health`: broker-4 disk: 87.3%; broker-4 leaders: 412/1680; clusterDiskAvg: 54.1%; underReplicated: 0
11. **Agent:** Found the root cause. Partitions 3 and 7 are both led by broker-4, which is at 87% disk utilization (cluster average is 54%). High disk pressure causes I/O contention that slows down consumer fetches. Broker-4 also holds 24.5% of all cluster leaders - well above the expected 20%.
12. **Operator:** what can we do about this?
13. **Agent:** Two options. A quick fix: trigger leader elections to move P3 and P7 off broker-4 immediately - zero data movement, instant effect on consumer lag. A lasting fix: generate a full rebalance proposal that redistributes both leaders and replicas to reduce broker-4's disk from 87% to a safe level. I'd recommend the full proposal since it prevents the problem from recurring. Want me to generate one?
14. **Operator:** yes, generate a proposal
15. **Tool call** `generate_proposal`: Running multi-phase optimization...
16. **Tool result** `generate_proposal`: moves: 8; leaderElections: 5; replicaTransfers: 3; dataToMove: 14.2 GB; diskVariance: -41%
17. **Agent:** Proposal ready: 8 moves total. 5 leader elections (instant, zero data transfer) and 3 replica moves (14.2 GB). Broker-4's disk drops from 87% to 71%, and its leader share of payments-events goes from 48% to 16%. The leader elections alone will immediately fix the consumer lag.
18. **Operator:** apply it
19. **Tool call** `apply_proposal`: Applying proposal #217...
20. **Tool result** `apply_proposal`: status: pending_approval; proposalId: 217; action: Requires human approval
21. **Agent:** This needs your explicit approval before executing. The 3 replica transfers will move 14.2 GB at the configured throttle rate. Leader elections are instant and zero-impact.
22. **Operator:** approved
23. **Tool call** `approve_action`: Executing approved proposal #217...
24. **Tool result** `approve_action`: status: executing; leaderElections: 5/5 complete; replicaTransfers: 0/3 in progress; throttle: 50 MB/s
25. **Agent:** Rebalancing started. All 5 leader elections completed instantly - partitions 3 and 7 now have new leaders on brokers with lower disk pressure. The 3 replica transfers are running at 50 MB/s, estimated ~5 minutes. Consumer lag should start dropping immediately.

Five messages from the operator. Seven tool calls by the agent. One approval gate. The entire investigation-to-fix loop completed in under a minute - a workflow that traditionally takes an hour or more of manual CLI work, JSON wrangling, and spreadsheet analysis.

## What Just Happened?

Let’s break down the agent’s investigation.

### 1. Consumer Group Inspection

The operator reported a symptom - slow order processing. The agent immediately called two tools in parallel: `describe_consumer_group` for group state and membership, and `get_consumer_group_lag` for per-partition lag numbers. Within seconds, it identified that lag was concentrated on two specific partitions (P3 and P7), not spread evenly - a pattern that points to an infrastructure issue rather than a slow consumer.

### 2. Root Cause Analysis

Without being asked, the agent followed the thread: if two partitions are lagging while others aren’t, what do those partitions have in common? It called `get_partition_activity` and `get_cluster_health` in parallel, and found the answer. Both partitions were led by broker-4, which was at 87% disk utilization while the cluster average was 54%. High disk pressure causes I/O contention that directly impacts consumer fetch latency.

This kind of multi-hop reasoning - from consumer lag to partition assignment to broker health - is exactly what makes Kafka troubleshooting time-consuming. The agent connected the dots in seconds.

### 3. Options, Not Orders

When the operator asked “what can we do?”, the agent didn’t jump straight to a solution. It presented two options: a quick fix (leader elections only - instant relief but temporary) and a lasting fix (full rebalance proposal that also addresses the disk imbalance). It explained the trade-offs and recommended the full proposal, but left the decision to the operator.

Once asked, the `generate_proposal` tool triggered Pilot’s multi-phase rebalancing pipeline. This isn’t a cluster-wide shuffle. The engine generated a targeted plan: 5 leader elections (instant, zero data transfer) to immediately redistribute read load, plus 3 replica moves (14.2 GB) to prevent the problem from recurring. Broker-4’s disk usage drops from 87% to 71%.

The distinction between leader elections and replica moves matters. Leader elections change which broker serves reads for a partition - they complete in milliseconds with zero data movement. Replica moves physically transfer partition data between brokers - they take time but fix the underlying imbalance. The engine maximizes elections and minimizes transfers.

### 4. Approval Gate

When the agent called `apply_proposal`, Pilot didn’t just execute it. The tool returned `pending_approval` and the agent explained what would happen: which moves are instant, how much data will transfer, and at what rate. Every mutation - whether it’s a partition reassignment, a config change, or an ACL update - goes through this gate. The AI can diagnose and recommend, but it cannot act without your consent.

## 60+ Tools, One Natural Language Interface

The MCP integration exposes Pilot’s full operational surface as structured tools that the AI can compose dynamically:

**Diagnostics** - `get_cluster_overview`, `get_cluster_health`, `get_partition_activity`, `get_broker_racks`, `search_topics`, `list_consumer_groups`, `describe_consumer_group`, `get_consumer_group_lag`

**Simulation** - `what_if_simulate` (broker failure, rack failure, traffic spike), `blast_radius_analyze` (single-failure impact analysis), `generate_proposal` (multi-metric optimization)

**Operations** - `apply_proposal`, `redistribute_topic`, `redistribute_partition`, `preferred_leader_election`, `update_topic_config`, `reset_consumer_group_offsets`, `create_acls`, `update_quota`

**Audit** - `search_audit_log`, `revert_audit_event` (undo previous changes)

Read-only tools execute automatically. Mutation tools always require approval. This isn’t a configuration option - it’s a hard architectural constraint.

## Why This Matters

Kafka operations have a knowledge problem. The platform is powerful but the operational surface is vast: broker configs, topic configs, partition assignments, consumer groups, quotas, ACLs, rack awareness, replication factors, ISR management. Most teams have one or two people who understand it deeply, and everyone else is searching Stack Overflow during incidents.

An MCP-powered agent doesn’t replace expertise - it makes expertise accessible. A junior engineer can report “orders are slow” and get the same root cause analysis that a senior SRE would produce. The approval gate ensures they can’t accidentally break anything.

The combination of natural language understanding, structured tool access, and human-in-the-loop safety creates an operations model where:

- **Diagnosis is instant** - trace symptoms to root causes across multiple subsystems
- **Risk is quantified** - understand exactly what a change will do before approving
- **Actions are safe** - approval gates on every mutation
- **Everything is auditable** - full trail of what was asked, what was done, and by whom

## Try It

Pilot is available as a Docker image. The AI assistant works with Anthropic (Claude), OpenAI, and local models - bring your own API key or run fully offline. Check out the [documentation](https://docs.calinora.io) for setup instructions, or visit [calinora.io](https://calinora.io) to learn more.
