Skip to Content
FeaturesArchitecture

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

Core Components

Kafka Client

A single connection pool handles every Kafka interaction:

SurfaceExamples
Metadata + healthCluster info, broker availability, partition watermarks, log directory sizes
ReadsMessage browser, audit consumer, state topic consumers
MutationsReassignments, 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):

  1. Watermark deltas - high watermark changes over time give message rates per partition
  2. Log directory sizes - disk usage per partition per broker gives byte rates
  3. Follower lag aggregation - the OffsetLag field from the DescribeLogDirs response is extracted and aggregated at cluster, broker, topic, and partition granularity, piggybacking on the same log directory call with no additional Kafka requests
  4. 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:

TopicPurposePolicyCreated
__pilot_broker_stateBroker status (online/offline/maintenance)CompactedAlways on startup
__pilot_audit_logAudit trail of all mutationsDelete (30-day retention)On startup when AUDIT_ENABLED=true (default)
__pilot_pat_tokensPersonal access token hashesCompactedWhen 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:

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:

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:

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.

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.

Last updated on