Curated summary
Husky: Exactly-once ingestion and multi-tenancy at scale
Husky, Datadog’s distributed, time-series-oriented event store, is optimized for large scans and aggregations rather than high-volume, low-latency point lookups. This makes exactly-once ingestion challenging, especially at Datadog’s multi-tenant scale. Datadog addresses the problem with deterministic, locality-aware routing that limits deduplication scope, improves storage efficiency, and supports autoscaling ingestion pipelines.
Husky’s Ingestion Challenge
- Husky separates storage and compute, allowing each to scale independently.
- Its storage engine is designed primarily for large analytical scans and aggregations.
- It is not optimized for massive numbers of low-latency point lookups, complicating duplicate detection during ingestion.
- The ingestion system must guarantee that every event is stored exactly once while maintaining:
- Multi-tenant scalability
- Reasonable ingestion latency
- Controlled infrastructure and storage costs
Routing Events to Storage Shards
- Datadog uses an upstream Shard Router to introduce locality into Kafka pipelines.
- Events are deterministically assigned to shards based on their tenant, timestamp, and event ID.
- Each tenant receives a list of shards rather than being permanently assigned to one shard.
- A deterministic choice from that list distributes the tenant’s events while keeping the number of active shards as small as practical.
- Downstream workers consume one or more shards and perform exactly-once ingestion into Husky.
Benefits of Data Locality
Simpler deduplication
- An event with the same timestamp and ID always reaches the same shard.
- Deduplication only needs to occur within that shard.
- Workers handle smaller sets of event IDs, making in-memory deduplication more efficient.
Lower storage costs and better performance
- Each shard processes a relatively small set of tenants.
- Husky stores each tenant in a separate table and does not mix tenants within files.
- More tenants per writer produce more output files, increasing blob-storage costs and compaction work.
- Restricting tenant cardinality reduces file creation and improves writer and compactor efficiency.
Challenges in Deterministic Routing
Changing shard assignments
- Tenant traffic can increase by one or two orders of magnitude.
- Assignments may change when scaling a tenant across more shards or rebalancing traffic among existing shards.
Distributed router consensus
- Every Shard Router node must make the same routing decision for a given event.
- Inconsistent decisions could send duplicates to different shards and undermine exactly-once ingestion.
Load balancing
- Shards must receive roughly equal traffic so downstream ingestion workers remain balanced.
Time-Bounded Shard Placements
A simple deterministic mapping can select a shard using a hash of the event ID:
shard = shards[hash(event_id) % num_shards]This approach is cheap and stateless when all routers know the same shard list.
However, changing the shard list can cause the same event ID to map to a different shard, so assignment changes require additional coordination.
The article introduces time-bounded Shard Placements to preserve consistent routing while allowing tenant assignments to evolve, though the supplied excerpt ends before explaining the mechanism in detail.
Datadog’s core recommendation is to combine deterministic, tenant-aware routing with carefully coordinated assignment changes. This narrows the scope of deduplication while reducing storage overhead and enabling balanced, scalable exactly-once ingestion.
Related reading
Continue with another curated summary.
Inside Husky’s query engine: Real-time access to 100 trillion events
Read originalAchieving relentless Kafka reliability at scale with the Streaming Platform
Read originalHow we measure data completeness at scale
Read originalHow we migrated a live routing system using AI-assisted refactoring
Read original