Skip to content

Latest commit

 

History

History
91 lines (53 loc) · 3.33 KB

File metadata and controls

91 lines (53 loc) · 3.33 KB

Phase 2+ Design Notes

Topic

Static vs Dynamic Partitions

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.

Key Routing

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.


Consumer Groups

Rebalancing

Trigger: heartbeat timeout (consumer dies), new consumer joins, partition count changes.

Process:

  1. Pause all consumers in the group
  2. Coordinator reassigns partitions (round-robin or range)
  3. Each consumer resumes from its last committed offset

Cost: stop-the-world pause. Kafka added cooperative rebalancing (incremental) — only affected partitions stop.

Coordinator Requirements

  • 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)

Offset Tracking

Offsets are per-partition. Consumer groups need per-partition-per-group offset storage. New partitions start at offset 0 — no effect on existing offsets.


Replication & Durability

fdatasync Window

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.

Recovery Truncation

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.


File Descriptor Management

Steady-state: 3 fds per partition (active segment only). Closed segments released after recovery (Segment.Release). Read opens files on demand.


Debugging Tools (Phase 2)

Protobuf Visibility

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.


Recovery Functions

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.