Static (Phase 1): Created upfront, count fixed. Simpler, matches Kafka default.
Dynamic (Phase 2+): Add partitions after creation. Use case: scaling throughput (more partitions = more consumer parallelism). Refactor is moderate — extract partition creation into AddPartitions(n) method.
Dynamic downside: hash(key) % N changes when N changes — same key routes to different partition, breaking per-key ordering. Requires consumer rebalancing.
Modulo (hash(key) % partitionCount): simple, uniform, all keys reassign on partition change. Kafka's choice.
Consistent hashing (ring hash): minimal key reassignment (~1/N keys move), but adds complexity (virtual nodes for even distribution) and doesn't solve consumer rebalancing.
Decision: modulo. Simpler, and rebalancing is needed for consumer groups anyway.
Trigger: heartbeat timeout (consumer dies), new consumer joins, partition count changes.
Process:
- Pause all consumers in the group
- Coordinator reassigns partitions (round-robin or range)
- Each consumer resumes from its last committed offset
Cost: stop-the-world pause. Kafka added cooperative rebalancing (incremental) — only affected partitions stop.
- Heartbeat protocol (detect dead consumers)
- Member list (which consumers are alive)
- Partition assignment strategy (round-robin, range, sticky)
- Offset storage per consumer group (per-partition-per-group committed offsets)
Offsets are per-partition. Consumer groups need per-partition-per-group offset storage. New partitions start at offset 0 — no effect on existing offsets.
Crash before fsync = recent writes in page cache are lost. Mitigated by replication: with ISR (min.insync.replicas=2), data survives on other replicas even if one crashes before fsync.
CRC failure, offset gap, or partial message stops scanning — remaining valid messages after corruption are lost. Can't trust read position after corruption (shifted bytes could make subsequent CRC checks pass on wrong data). Kafka truncates at corruption too.
Steady-state: 3 fds per partition (active segment only). Closed segments released after recovery (Segment.Release). Read opens files on demand.
2PC adds Prepare/Commit/Abort message flows. Need visibility during development.
Options:
- Server flag
-log-protobuf— dumps decoded requests/responses to stderr - Standalone TCP proxy (~100 lines Go) — sits between client and server, decodes frames, logs, forwards
protoc --decode— offline decoding of captured traffic
Decision: build server flag + proxy during Phase 2. Not now.
| Function | Purpose |
|---|---|
RecoverSegment |
Open + scan + truncate index for a single segment |
recoverActiveSegment |
Same as RecoverSegment but O_RDWR for writes |
Recover |
Orchestrate recovery across all segments in a partition |
Release |
Close fds without syncing (closed segments after recovery) |
Index/timeindex functions (openIndex, openTimeIndex) just open files — no recovery logic.