Architecture
Pilot is a single binary that connects to your Kafka cluster, continuously samples partition metrics, generates rebalancing proposals, and executes reassignments - all without an external database.
High-Level Data Flow
graph TD
K["Kafka Cluster<br/><small>Brokers · Topics · Partitions</small>"]
K -->|"metadata + health"| CC["Kafka Client<br/><small>single connection pool</small>"]
K -->|"watermarks + disk sizes"| MS["Metadata Sampler<br/><small>activity scoring</small>"]
CC -->|"broker health<br/>partition metadata"| PE["Proposal Engine<br/><small>multi-phase pipeline</small>"]
MS -->|"7 metrics per broker"| PE
PE -->|"rebalancing proposal"| RE["Reassignment Executor<br/><small>streaming submission · throttle mgmt</small>"]
RE -->|"alter_replica_assignment"| K
Core Components
Kafka Client
A single connection pool handles every Kafka interaction:
| Surface | Examples |
|---|---|
| Metadata + health | Cluster info, broker availability, partition watermarks, log directory sizes |
| Reads | Message browser, audit consumer, state topic consumers |
| Mutations | Reassignments, config updates, leader elections, ACLs, quotas |
| State persistence | __pilot_broker_state, __pilot_audit_log, __pilot_pat_tokens |
Metadata Sampler
Pilot infers partition activity from Kafka metadata APIs - the sampling loop does not consume topic messages.
Every collection cycle (default 10s):
- Watermark deltas - high watermark changes over time give message rates per partition
- Log directory sizes - disk usage per partition per broker gives byte rates
- Follower lag aggregation - the
OffsetLagfield from theDescribeLogDirsresponse is extracted and aggregated at cluster, broker, topic, and partition granularity, piggybacking on the same log directory call with no additional Kafka requests - Follower calibration - learns how much replication traffic followers actually generate relative to leader writes, smoothed over time per topic class
These are combined into seven metrics per broker: leader count, follower count, disk bytes, producer message rate, producer byte rate, consumer message rate, and consumer byte rate.
The activity score is a composite indicator that blends message rate, byte rate, disk size, and cumulative volume into a single 0-100 number per partition. Each component is log-normalized so large traffic differences don’t dominate the score. It is a dashboard “how busy is this broker?” summary only - the rebalancing optimizer works directly off the seven metrics above. The activity score is deliberately not one of them: it is a capped re-encoding of metrics the optimizer already balances, so including it would double-count producer throughput and disk in the objective.
State Persistence
Pilot stores state in Kafka itself - no external database:
| Topic | Purpose | Policy | Created |
|---|---|---|---|
__pilot_broker_state | Broker status (online/offline/maintenance) | Compacted | Always on startup |
__pilot_audit_log | Audit trail of all mutations | Delete (30-day retention) | On startup when AUDIT_ENABLED=true (default) |
__pilot_pat_tokens | Personal access token hashes | Compacted | When PATs are enabled |
__pilot_broker_state and __pilot_pat_tokens are created with 1 partition and the cluster’s default replication factor (1 if that fails). The audit topic follows AUDIT_PARTITIONS and AUDIT_REPLICATION_FACTOR; see Audit Logging. The cleanup policy is enforced idempotently on every startup, so manual topic misconfiguration is self-correcting.
See Audit Logging for audit topic configuration and Authentication for PAT configuration.
Proposal Pipeline
The proposal engine generates rebalancing plans through a multi-phase pipeline:
graph TD
CS["Cluster State"] --> PREP
subgraph PREP["Preparation"]
direction TB
A1["Fetch metadata & check safety"] --> A2["Filter exclusions & fix critical issues"]
end
PREP --> OPT
subgraph OPT["Optimization"]
direction TB
B1["Multi-level search"] --> B2["Cost reduction"]
end
OPT --> FIN
subgraph FIN["Finalization"]
direction TB
C1["Spread & validate"] --> C2["Prune moves & rebuild metrics"] --> C3["Check benefit & diagnose limits"]
end
FIN --> OUT["Proposal"]
Preparation collects metadata, checks cluster safety (in-flight reassignments, ISR shrinks), filters excluded or cooling-down partitions, and fixes critical issues like under-replicated partitions and rack violations. Brokers with active follower lag are marked so the optimizer avoids sending them new inter-broker data. If every broker is within 10% of its fair share of every load with traffic (see No Traffic) and no topic is concentrated, the pipeline short-circuits here.
Optimization is the core of the engine. It searches for partition moves that reduce imbalance using focused, broad, and swap candidate strategies. Each candidate is scored by how much it improves balance versus how expensive the move is. A second pass then reduces proposal cost by pruning unnecessary moves and replacing expensive inter-broker transfers with cheaper leader-only shifts.
Finalization moves leadership off concentrated topics, validates correctness, and evaluates each partition’s net transition. It removes no-ops and changes that do not help, then rebuilds projected broker metrics from the final placements. Optional plans must still be worth applying after pruning and any partition cap. If correctness repairs are needed, or a broker is down and not in maintenance, the final proposal contains only repairs and defers optional balancing. Loads at their limit produce advisories that say what would help.
Completed transitions start a 30-minute cooldown, in which optional balancing leaves those partitions alone. See Reducing Unnecessary Movement for the worth-applying check and Recent Moves for the cooldown.
Application admission checks measurement readiness and evaluates the exact final plan against three fresh, compatible observations from the bounded measurement history. Separate optimizer runs do not need to select matching moves. Manual and automatic submission recheck the plan against current inputs; changes to its starting placements or destination eligibility require recalculation.
Scoring
Metrics are weighted by operational priority. Leader and follower counts have the highest weight, followed by disk usage, then byte rates, then message rates.
Movement Cost Model
The engine prefers cheaper moves over expensive ones. Leader-only changes (no data transfer) are cheapest, followed by inter-broker moves (proportional to partition size), with cross-rack transfers being the most expensive due to network cost.
Reassignment Execution
When a proposal is applied, the executor uses streaming submission:
graph TD
P["Proposal<br/><small>N partition moves</small>"] --> Q["Pending Queue"]
Q --> B1["Broker 1<br/><small>active: 3/5 → submit</small>"]
Q --> B2["Broker 2<br/><small>active: 5/5 → wait</small>"]
Q --> B3["Broker 3<br/><small>active: 1/5 → submit</small>"]
B1 --> KA["Kafka<br/><small>alter_replica_assignment</small>"]
B3 --> KA
KA --> RM["Reassignment Monitor<br/><small>Track progress · Detect stuck · Retry</small>"]
RM -->|"completed"| Q
Unlike batch/wave execution, streaming continuously submits new moves as soon as a broker has capacity. This maximizes cluster utilization while respecting per-broker concurrency limits (PILOT_MOVE_MAX_PER_BROKER).
Both the concurrency limit and throttle rate (PILOT_THROTTLE_RATE_MB) can be adjusted live while a reassignment is in progress - changes take effect immediately without restarting. Throttle configs are cleaned up automatically after reassignment completes.
Self-Healing
Three independent background loops run at different intervals:
graph LR
subgraph Healing Loops
direction TB
C["Critical Fixes<br/><small>every 5m · highest priority</small><br/><small>URP repairs, rack violations</small>"]
R["RF Increases<br/><small>every 15m · medium priority</small><br/><small>Replication factor adjustments</small>"]
A["Activity Balancing<br/><small>every 30m · lowest priority</small><br/><small>Load rebalancing across brokers</small>"]
end
C -->|"blocks"| R
R -->|"blocks"| A
C --> G{Safe to apply?}
R --> G
A --> G
G -->|"yes"| E["Execute Proposal"]
G -->|"no"| S["Skip<br/><small>cooldown / window / dry-run</small>"]
Each loop obtains a proposal, checks safety signals, and applies eligible changes. Critical fixes run first, then RF increases, then activity balancing. While one loop runs, the others skip their cycle. All loops respect time windows, cooldowns, and dry-run mode. Activity balancing shares the default manual application’s measurement readiness, sustained-benefit and freshness requirements. Its per-run partition limit is included in generation so the final plan’s benefit is checked before application.
AI Integration
Chat Assistant
The chat engine connects a large language model to Pilot’s full API surface. Supported providers are Anthropic and OpenAI. You can also point at self-hosted or third-party models via PILOT_CHAT_BASE_URL, as long as they implement the Anthropic or OpenAI chat API format.
graph LR
U["User Message"] --> LLM["LLM<br/><small>Claude / GPT</small>"]
LLM -->|"tool call"| MCP["MCP Tool Executor"]
MCP --> RO{Read-only?}
RO -->|"yes"| EX["Execute"]
RO -->|"no"| AP["Approval Request"]
AP -->|"approved"| EX
AP -->|"denied"| DN["Return denial"]
EX -->|"tool result"| LLM
DN -->|"tool result"| LLM
Mutation tools (reassign partitions, alter configs, delete ACLs) require explicit user approval before execution. Read-only tools execute immediately.
Model Context Protocol (MCP)
The MCP server exposes 56 tools that external AI agents can invoke over HTTP. Tools are categorized as read (auto-execute) or mutate (require approval). See MCP for the full tool catalog.