OpenTelemetry at Scale: Using Kafka to handle bursty traffic
signoz.io
signoz.io
Grafana has some posts how they softened the s3 blow with memcached(2,3).
1. https://github.com/open-telemetry/opentelemetry-collector-co... 2. https://grafana.com/docs/loki/latest/operations/caching/ 3. https://grafana.com/blog/2023/08/23/how-we-scaled-grafana-cl...
I know the post is about telemetry data and my comments on grafana are logs, but the arch bits still apply.
The linked Loki caching docs/articles are for optimising the read access patterns of S3/object storage, not for writes.
One thing that is also very interesting with Kafka is that you can achieve exactly-once semantic without too much efforts: by keeping track of the positions of partitions in your own database and carefully acknowledging them when you are sure data is safely stored in your db. That's what we did with our engine Quickwit, so far it's the most efficient way to index data in it.
One obvious drawback with Kafka is that it's one more piece to maintain... and it's not a small one.
Some folks in SigNoz community have also suggested NATS for this, but I have not deep dived into benchmarks/features yet
That's not really "exactly once". What happens when your system dies after it made sure the data is safely stored in the db and before ack-ing?
What if the application doesn't restart before the queue decides the message was lost and resends?
One-or-more semantics + local deduplication gives one-and-only semantics.
In this case you're optimising local deduplication with strictly monotonic index.
One downside is that you leak internals of other system (partitions).
The other is that it implies serialised processing - you can't process anything in parallel as you have single index threshold that defines what has been and what has yet not been processed.
edit: You added some more to your comment after I posted this one, so I'll try to cover them as well:
> One downside is that you leak internals of other system (partitions).
Yeah, sure.
> The other is that it implies serialised processing - you can't process anything in parallel as you have single index threshold that defines what has been and what has yet not been processed.
It doesn't imply serialised processing. It depends on the use-case, if each record in a topic has to be processed serially, you can't parallelize full-stop; number of partitions equals 1. But if each record can be individually processed you get parallelism equal to the number of partitions the topic has configured. You also achieve parallelism in the same way if only some records in a topic needs to be processed serially, at which point you can use the same key for the records needing to be serially processed and they will end up in the same partition, for example recording the coordinates of a plane - each plane can be processed in parallel, but an individual plane's coordinates need to be processed serially - just use the planes unique identifier as key and the coordinates for the same plane will be appended to the log of the same partition.
If one-and-only-one semantics are needed and processing should be parallel, other methods have to be used.
> One downside is that you leak internals of other system (partitions).
True, but we generalized the concept of partitions for other datasources, pretty convenient to use it for distributing indexing tasks.
Fortunately Kafka is partitioned. You cannot work in parallel along partitions.
Also, you can streamline your process. If you are running your data through operation (A, B, C). (C on batch N) can run at the same time as (B on batch N+1), and (A on batch N+2)
We do both at quickwit.
You can make the downstream idemptoent wrt what the queue is delivering, but the queue might still redeliver things.
One simply trades latency for capacity and eventual coherent data locality.
Its almost a arbitrary detail whether you use Kafka, RabbitMQ, or Erlang channels. If you can add smart client application-layer predictive load-balancing, than it is possible to cut burst traffic loads by a magnitude or two. Cost optimized Dynamic host scaling is not always a solution that solves every problem.
Good luck out there =)
[^1] https://x.com/donkersgood/status/1662074303456636929?s=20
The point is to only use s3 etc in the event of system instability. Not as a primary data transfer means.
At that point it's very easy to sleepwalk into implementing your own database on top of s3, which is very hard to get good semantics out of - e.g. it offers essentially no ordering guarantees, and forget atomicity. For telemetry you might well be ok with fuzzy data, but if you want exact traces every time then Kafka could make sense.
The approach here is to only send data to s3 as a last ditch resort.
This is more or less exactly what WarpStream is: https://www.warpstream.com/blog/minimizing-s3-api-costs-with...
Kafka API, S3 costs and ease of use
Early days, I looked up otel and observability stuff, and I always saw Signoz articles on the first screen.
I'm curious how long things stay in Kafka on average and worse case. If it's more than a few minutes, I imagine it lowers the quality of tail based sampling.