Rust AI & Data Meetup: Kranti Parisa, building a fast streaming service in Rust with Apache Iggy
Uh for the intros I'm granted by Lisa and the founder CEO of laser data and also we are the co-creators. We are the creators of Apache game the co-creator of it. Um so why why did we you know went into this journey? Like Shahab mentioned about Spark to sale so Kafka to Iggy. So um So yeah, like taking a step back this is this is Rust and AI why latency matters for us for AI use cases. So this is a simple use case out of my own experience when we are building the voice bots. We ran into issues of like tail latencies were actually killing the agents really badly, right? So when when everything looking at P50 and averages everything looks good when you are operating within a Java based systems but then you know the tail latencies are the ones which are breaking the you know user experience, right? So Um the usual way of solving the problem is throwing hardware at it, right? You know I've done that you know many many times when I was working at Apple. I led search and personalization at Apple. And then like that's where the first love with Rust came into picture because we rewrote one of the spelling modules in the search pipeline and then the latencies were like significantly better
Like you know from P99, P95 latencies of like 100 to 200 milliseconds to 2 to 3 milliseconds. It's insanely fast, right? So that's that's when you know that the first love with Rust started. It was always there in my in my mind about like you know okay what about the event log event streams. So why does that really matter? Why why event streams and logs? Why did we pick up that problem to go solve out of like you know many things? So in the AI world when we were building these you know multi-agent agentic systems when I was leading engineering at Dialpad, we were building in voice bots and chatbots and things like that. And when when we were building multi-agents this is the N by M mesh kind of architecture, right? Where you have the harness system, you have the orchestrator, the request going into the orchestrator, and the orchestrator now calling the LLM you know a bunch of tools and like you know making MCP calls and then you know deciding which tool to hit and then it has to log into the observability and then it has to send the same data into your analytics pipeline and so on and so forth. So the orchestrator was becoming a bottleneck for us. So which because it is doing so much work. So many MCP calls, so many tool calls
Of course, we're not talking about replacing MCP and you know replacing tool calls that that has to be there for the external systems. But within the infrastructure itself, we don't need to use MCP because of you know the data exchange between sub-agents and the subsystems in your architecture can happen through low-level protocols like TCP or you know quick and so on and so forth. So we were like okay, we have seen this kind of pattern in the past, right? So that's how distributed systems have evolved you know for for two decades where we all started off with like oh internet like oh rest APIs. Oh great like you know just make calls to you know everywhere, right? But then we realized that like you know this is some of the systems needs to be asynchronous. Some of the systems need to be even driven. Right? And when when we are making MCP calls, we assume that the agent that we are calling is live and up and running, but it may not be the case in all the time. When when that goes down then everything you know is kind of like lost in the middle because it's not a persistent protocol, right? So then we started like you know okay, let's put a stream in between so that we make all these, you know, sub agents subscribe to that whenever there is a context change. Say, for example, a request is, "Refund my uh order number 1 2 3." Let's say if some, you know, customer support use case come in and then they ask for it
The agent, the orchestrator, now initiated this request. Okay, check whether this customer is, you know, legitimate or not. Check the payment was processed or not. Check like, you know, what are the CRM records, you know, whether the item has shipped or not, all of that stuff. Then, while it is doing and blasting the compute, then the user basically said, "So, oops, sorry, it was not 1 2 3, it's 1 2 3 4." Now, the entire, you know, blast radius of the whole ecosystem needs to be notified immedi- as quickly as possible. That's where the latency matter a lot, uh mattered to us a lot, and that's how uh we we we said, "Okay, this is not scaling. This is not working for the traditional uh event, you know, streams like Kafka, for example, right?" So, then we started like, you know, thinking about, "What do we do?" And then, you know, going back, you know, connecting the dots that, you know, the Apple journey of like, you know, replacing JVM with Rust, and then we were like, "Okay, let's let's go do it." That's how the Iggy was born, uh the streaming engine, and then later I joined the community, and then we moved into the, you know, Apache Foundation. So, okay, so going, you know, into the internals, like, you know, probably not going to go too much deep into it
Beyond Rust and beyond the just language itself, we have fine-tuned a lot of things because there is so much evolved in the last, you know, decade plus since, you know, Kafka and the legacy streaming engines were born, uh disks have changed, right? Like, the hardware has changed. The compute has changed, you know, from number of cores that are available. And then, you know, your IO, uh how you interact with IO has changed. And then, obviously, [clears throat] like, you know, Java to like, you know, no GC and no pauses, and you know, Rust have evolved as well. So, we wanted to take advantage of all of them, not just like you know just rebuild uh you know an existing pattern with you know a new programming language. Like you know it's not a wipe coded thing as such. It is really to Alexy's point it is engineering. So why you know we we we said it's not right itself writing itself is not the main main challenge
Our our argument for for the streaming it it is how you exchange data. It is how you like you know exchange metadata and so on and so forth. And I'm not going to go into too much of details because in the interest of time, but there are a bunch of blog posts that we have written on iggy.apache.org. Please go check it out. Like it it it it tells you how we have handled direct IO. How did we eliminated page cache? How did we optimized for IO Uring? Like you know how what happens when you do thread per core architecture where if you just you know issue a request to to the compute, then it kind of like you know jumps around between the CPUs, right? So instead of that when we pin a thread per core magic happens, right? So that's the thread per thread per core architecture. So we did a lot of those things and and it solved complex problems and it solved the tail latency. I'll show you like you know what what I really mean by that
There's it's very easy to get start with. It's just like you know cargo add and then like you know it's simple interface. We do have support for web sockets also because one of the AI use cases is how do you interact with the server to the client? So we built a web socket protocol as well. Kafka doesn't have it. And then we also have like SDKs in all the major languages. So clustering is in in beta mode right now. And then there is like Kafka proxy. There's so many people are asking about can we have that drop in replacement like you know what sale did for you know for spark
We intentionally didn't want to do it to begin with because we don't want it to carry the baggage of the like 15 years old protocol. But then you know, somebody in the community because we are in Apache Foundation, so we we can't say no for the community asked, right? So, they're building a proxy, but it's still core inside and inside the Rust engine. So, I'm talking about, okay, what does that give us? Like, you know, now we're talking about the tail latency. So, these numbers are not averages and P90 and P95 you observe. These are P99. And the beauty with Rust-based engines is like the tail latencies are really flat, right? Where you won't actually see a huge spike between P50, P90, and P99, right? Like they'll be like maybe 1 ms difference between between them, right? That's That's what we want for real-time AI use cases, right? So, that enabled us to go and build a commercial company on on top of this open source. So, that's what you know, we're building with Laser Data Cloud. So, Laser Data Cloud, you know, can give you two products or offerings
One is the managed service of EG, Apache EG itself, like, you know, very similar to Confluent, what Confluent did for Kafka. You know, you can run in any cloud, AWS, and GCP, and all of that stuff. But we didn't want to stop there because we want to go and solve the actual AI use case and AI AI problem. Then we started thinking about, okay, what does AI agents need? They need all these primitives of memory and context and state and so on and so forth. And going back to the problem that I was talking about, orchestrator was doing, "Oh, the memory is in vector DB." And then, you know, the context is in Redis store. And then something else in somewhere else and all of that stuff. So, we were like, "Okay, we don't need to do all of that stuff. We can just use one SDK
We can write into the stream and configure the stream to be able to powered by materialized views and, you know, vectors and, you know, all of that stuff." So, it's one and one one SDK that gives you all of those primitives. That's what Laser Data the Cloud offers. And and as I said, like, you know, we we have bunch of these primitives. It's pretty simple language, Uh, simple SDK. We have the SDK in TypeScript and Python and Rust. Uh, please give it a try and you know, let us know uh, how it goes. THANK YOU. >> [applause] >> ANY QUESTIONS? YEAH
JUST GIVE ME YOUR FIRST. YOU can begin. >> Hi. Oh, thank you for the talk. So, I wouldn't say I'm an expert on this, but from my understanding, like Kafka comes with like a lot of overhead, right? And so, for laser data in particular, at what point does it make sense for you to implement this work kind of an orchestration layer, I would say. But, at what point at what scale does it make sense to kind of use something like Apache Iggy? >> So, Apache Iggy, uh, right now it's a it's a standalone single node, uh, engine, but we are working on the clustering mode. The minute that, you know, clustering goes out, which is probably, you know, a couple of weeks from now. Uh, and by the way, we haven't implemented raft protocols and you know, those sorts of protocols for Zookeeper-based, you know, uh, consensus and all of that stuff
So, we took an MIT paper. Uh, it's called view stamp replication. So, we have built a replication engine which is native. Uh, it's a single binary. Uh, so yeah, like any workload, like you know, if you want consistent tail latencies and reduction of cost, like you know, the cost reduction is really like 2/3 of the reduction in the cost and like 30x, you know, faster uh, tail latencies. So, so those are the like, you know, use cases. And then the orchestration layer that you're talking about, that is the the second uh, aspect of it, right? Iggy, Apache Iggy itself is the core streaming engine. You can just use it for general-purpose distributed systems, event-driven architectures, just like Kafka, you know, what Kafka did, but probably more because we do support web sockets, so you can actually talk to the clients also, not just the between the server to server, but server to clients also
Like the end end, you know, clients' devices and all of that stuff. But then, you know, the orchestration stuff is the Laser SDK, which is the AI native, you know, uh use cases where you want a persistent layer between agents. And the beauty of, you know, streaming having a stream or a log in between is you can configure all of your secondary services behind it, right? Like all of your observability, analytics, search, and all of the use cases, the data the agents are exchanging between them through the log will be automatically dumped down into the secondary services. So, the orchestrator becomes lighter, thinner, and, you know, much easy to scale because that need doesn't need to do all of this other work of observability and, you know, what's going on and replayability and all of that stuff. If it makes sense. >> One more question. All right. So, um it is it's been my experience with um thread per core versus work stealing that thread per core really heavily depends on appropriate partitioning of workloads at the beginning
And I'm curious how you solved that. >> Um yeah, I mean, at the end of the day, for us, uh I think that's what I was trying to mention here also. So, when we think about log and append-only log, we have to go with partitions anyways, right? So, that's how we have pinned, you know, for each of these partitions, then like, you know, pin for uh for each one of them. That's how, you know, we took advantage of the the the pattern. OH, THANK YOU. >> [applause] >> HEY.