6 ms·
OpenTelemetry at Scale: Using Kafka to handle bursty traffic
- francoismassot 3y agoI heard several times that Kafka was put in front of elasticsearch clusters for handling traffic burst. You can also use Redpanda, Pulsar, NATS and other distributed queues. 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.
- pranay01 3y agoHave you done/seen any benchmarks between Redpanda/NATS and Kafka for this use case? Some folks in SigNoz community have also suggested NATS for this, but I have not deep dived into benchmarks/features yet
- francoismassot 3y agoUnfortunately no :/
- radicality 3y agoExactly-once semantics of what specifically? Or do you mean at-least-once ?
- francoismassot 3y agoExactly-once semantic between Kafka and the observability engine.
- richieartoul 3y agoYou have to do a bit more than that if you want exactly once end-to-end (I.E if Kafka itself can contain duplicates). One of my former colleagues did a good write up on how Husky does it: https://www.datadoghq.com/blog/engineering/husky-deep-dive/ https://www.datadoghq.com/blog/engineering/husky-deep-dive/
- francoismassot 3y agoYeah, I was only talking about exactly once semantic between Kafka and Quickwit.
- viraptor 3y ago> 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 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?
- Svenskunganka 3y agoDepending on how you use the database it is. If you write the data as well as the offset to the DB in the same transaction, you can then seek to the offset stored in the DB after application restart and continue from there.
- deleted 3y ago[deleted]
- viraptor 3y ago> after application restart and continue from there. What if the application doesn't restart before the queue decides the message was lost and resends?
- hashhar 3y agoIn Kafka the "queue" is dumb, it doesn't lose messages (it's an append only durable log) nor does it resend anything unless the consumer requests it.
- viraptor 3y agoThere has to be a retry system somewhere, otherwise you'd end up with a 0-or-more delivery system if the app crashes after picking up from the queue, but never processing or ack-ing.
- mirekrusin 3y agoYou should drop "(...) and carefully acknowledging them when you are sure data is safely stored in your db (...)" part then, because it means it's not necessary, you don't rely on it. 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.
- foota 3y agoIsn't exactly once delivery the kind of problem like the CAP thereom where it's not possible? You can make the downstream idemptoent wrt what the queue is delivering, but the queue might still redeliver things.
- ankitnayan 3y agohttps://www.confluent.io/blog/exactly-once-semantics-are-possible-heres-how-apache-kafka-does-it/ https://www.confluent.io/blog/exactly-once-semantics-are-pos...
- bushbaba 3y agoSeems like overkill no? Otel collectors are fairly cheap, why add expensive Kafka into the mix. If you need to buffer why not just dump to s3 or similar data store as a temporary storage array.
- francoismassot 3y agoI really like this idea. And there is an OTEL exporter to AWS S3, still in alpha but I'm gonna test it soon: https://github.com/open-telemetry/opentelemetry-collector-contrib/tree/main/exporter/awss3exporter https://github.com/open-telemetry/opentelemetry-collector-co...
- prpl 3y agoWhy not both, dump to S3 and write pointers to kafka for portable event-based ingestion (since everybody does messages a bit differently)
- bushbaba 3y agoNo need as s3 objects is your dead letter queue and the system should be designed anyway to coupe with multiple write of same event. The point is to only use s3 etc in the event of system instability. Not as a primary data transfer means.
- lmm 3y ago> If you need to buffer why not just dump to s3 or similar data store as a temporary storage array. 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.
- dikei 3y agoYeah, and to use S3 efficiently you also need to batch your messages into large blobs of at least 10s of MB, which further complicates the matter, especially if you don't want to lose those messages buffers.
- daurnimator 3y agoI expect it would be far cheaper to scale up tempo/loki than it would be to even run an idle kafka cluster. This feels like spending thousands of dollars to save tens of dollars.
- neetle 3y agoTempo can still buckle under huge bursts of traffic, and you don’t need the retention to be in the hours
- deleted 3y ago[deleted]
- pranay01 3y agoQuerying in Tempo/Loki does seem to not scale particularly well, and Loki has known issues with high cardinality data, so...
- ankitnayan 3y agoWhen handling surges of the order of 10x, it's much more difficult to scale the different components of loki than to write them to Kafka/Redpanda first and consume at a consistent rate.
- Spivak 3y agoWhere are you finding such an expensive Kafka cluster? Kafka can run on 3 VPS's in a trenchcoat.
- blinded 3y agoThis arch is how the big players do it at scale (ie. datadog, new relic - the second it passes their edge it lands in a kafka cluster). Also otel components lack rate limiting(1) meaning its super easy to overload your backend storage (s3). Grafana has some posts how they softened the s3 blow with memcached(2,3). 1. https://github.com/open-telemetry/opentelemetry-collector-contrib/issues/6908 https://github.com/open-telemetry/opentelemetry-collector-co... 2. https://grafana.com/docs/loki/latest/operations/caching/ https://grafana.com/docs/loki/latest/operations/caching/ 3. https://grafana.com/blog/2023/08/23/how-we-scaled-grafana-cloud-logs-memcached-cluster-to-50tb-and-improved-reliability/ 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.
- ankitnayan 3y agoCaching is to improve read performance whereas Kafka is used to handle ingest volume. I couldn't correlate the Grafana articles shared
- blinded 3y agoYep! Should have made that more clear. Brought it up as an example of other parts of the system that require scaling.
- wardb 3y agoGrafana Labs employee here => On the linked articles: I'm not aware of any caching being used in the writing data to S3 part of the pipeline other then some time based/volume based buffering at the ingester microservices before writing the chunks of data to object storage. The linked Loki caching docs/articles are for optimising the read access patterns of S3/object storage, not for writes.
- blinded 3y agoYes. Thanks for the reply!
- kilotaras 3y ago
- deleted 3y ago[deleted]
- chris_armstrong 3y agoA similar idea [^1] has cropped up in the serverless OpenTelemetry world to collate OpenTelemetry spans in a Kinesis stream before forwarding them to a third-party service for analysis, obviating the need for a separate collector, reducing forwarding latency and removing the cold-start overhead of the AWS Distribution for OpenTelemetry Lambda Layer. [^1] https://x.com/donkersgood/status/1662074303456636929?s=20 https://x.com/donkersgood/status/1662074303456636929?s=20
- nicognaw 3y agoSignoz is too good at SEO. Early days, I looked up otel and observability stuff, and I always saw Signoz articles on the first screen.
- relaxing 3y agoBizarre. There's so little technical detail in the post beyond "use Kafka to handle bursty traffic", which is like, duh.
- Joel_Mckay 3y agoIf you have distributed concurrent data streams that exhibit coherent temporal events, than at some point you pretty much have to implement a queuing balancer. 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 =)
- nijave 3y agoIt'd be nice to have something simpler as an otel processor. Otel could just dump events to local disk as sequential writes then read them back, load permitting. 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.
- anacrolix 3y agoAre there any client side dynamic samplers that can target a maximum event rate? Burstiness with otel has been a thorn in everything that uses it from my experience and it's frustrating.
- pranay01 3y agoDo you mean sampling at application level before sending traces to otel collector/Kafka?