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) . 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.
- Operator: the payments team says orders are processing slowly. can you check their consumer group?
- Tool call
describe_consumer_group: group: "payments-processor" - Tool call
get_consumer_group_lag: group: "payments-processor" - Tool result
describe_consumer_group: group: payments-processor; state: Stable; members: 3; topics: payments-events - Tool result
get_consumer_group_lag: totalLag: 847,231; partitions: 12/12 assigned; maxLag: P3: 312K, P7: 289K; avgLag: 70.6K - 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.
- Tool call
get_partition_activity: topic: "payments-events" - Tool call
get_cluster_health: Checking broker health... - Tool result
get_partition_activity: P3 leader: broker-4; P7 leader: broker-4; broker-4 leaders: 48% of topic; throughput: 1,240 msg/s - Tool result
get_cluster_health: broker-4 disk: 87.3%; broker-4 leaders: 412/1680; clusterDiskAvg: 54.1%; underReplicated: 0 - 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%.
- Operator: what can we do about this?
- 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?
- Operator: yes, generate a proposal
- Tool call
generate_proposal: Running multi-phase optimization... - Tool result
generate_proposal: moves: 8; leaderElections: 5; replicaTransfers: 3; dataToMove: 14.2 GB; diskVariance: -41% - 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.
- Operator: apply it
- Tool call
apply_proposal: Applying proposal #217... - Tool result
apply_proposal: status: pending_approval; proposalId: 217; action: Requires human approval - 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.
- Operator: approved
- Tool call
approve_action: Executing approved proposal #217... - Tool result
approve_action: status: executing; leaderElections: 5/5 complete; replicaTransfers: 0/3 in progress; throttle: 50 MB/s - 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 for setup instructions, or visit calinora.io to learn more.