SBTB 2023: Boyang Jerry Peng, Latency goes sub-second in Apache Spark Structured Streaming.
Hi everyone. My name is Jerry and I'm here to talk about Project Lightspeed and some of the the work we've been doing on improving the performance and functionality of Spark Streaming. Um the original presenter was going to be Karthik. Clearly, I'm not Karthik. Karthik is unfortunately feeling under the weather, so I am going to pres- present on his behalf. For those that were here yesterday, I was also part of the the panel to talk about streaming. So, let's get started. But, I am also from Databricks
I work with Karthik on streaming at Databricks. So, let me first talk about what is stream processing. So, stream processing, just to give like an overview, is basically you're processing continuous and unbounded amounts of data in real time. Um you're doing some sort of transformation like, you know, window aggregation, it could be pattern detection, enrichment, routing, you know, if you're more familiar with SQL or databases, could be projections, you know, uh filtering, you know, all of all of that stuff, you know, as well. And basically, you're processing data that's coming in live, you're doing some sort of transformation, and you're getting some sort of result. And there could be a broad range of use cases that I'll briefly, you know, um go over. So, for example, like, you know, in more recent years, just like what we discussed in the streaming panel yesterday, there has been an explosion of kind of use cases that require streaming and real-time processing. A lot of it comes from like, you know, the areas I've list that's are listed here
Financial services, one of the, you know, simple examples to reason about is fraud detection. Why that would require kind of real-time processing that if there are credit card transactions or transactions of any kind that um can be classified as fraudulent, you want to be able to detect those early and notify the, you know, responsible parties. Retail, you know, personalization. If you go to Amazon or go to any sort of online retail, you know, the retailer would like to be able to recommend things in real time to you to take advantage of the opportunities so that the consumer can be aware be aware of these products and, you know, purchase them. Healthcare, you know, we've been through COVID-19 and other um health-related uh issues and being able to kind of detect um basically uh trends in real time is generally helpful as well. Manufacturing, one of the biggest use cases is predictive maintenance. If, say, on an assembly line somewhere, you know, you you're able to detect that some component on the assembly line was able was going to fail or was going to be failing in the near future, you want to be able to proactively be able to replace those kind of um components to to minimize the amount disruption time uh that would uh incur on your kind of assembly line. Energy
In the energy kind of industry, smart pricing would be a good example of, you know, where for people living in California, you know, like different times of the day, there's different like, you know, pricing for, say, electricity. How like, say, PG&E or other companies, energy companies could basically detect the amount of energy used in the grid and be able to price it smartly to incentivize people to use different energy and, you know, different times, you know, would you know, be you know, beneficial to them. Gaming, there's a lot of examples there, you know, interactive analytics. And it also kind you know, a lot of the other things could be fraud detection as well or cheating in video games as well could be could benefit uh from real-time processing to detect these kind of instances in real time. Um technology and software for a lot of like, um you know, digi- digital native companies or, you know, companies like building smart cars, smart homes, and things like that. Lots of data coming in. How can I basically, you know, take advantage of opportunities or be able to send out notifications to inform users of some event kind of in real time um you know, in those areas. Media entertainment
Um again, this kind of is ties similar to retail is con- is recommendation, right? Can we provide real-time recommendations to end users uh for them to, you know, basically consume content that they're, you know, looking for. So, let me briefly talk about like how Spark Streaming or Structured Streaming Sorry, I didn't mean to my slide. Um is is has has been doing. So, this is actually kind of old data, so it actually has grown even more so than this, but there we have seen tremendous growth in just, you know, adoption and usage of Spark Structured Streaming, you know, year over year, especially, you know, after 2001 and the and even more recently, we've seen, you know, even a a bigger growth and adoption as we kind of Databricks has invested more in this area. So, some of the customers, um we have quite a few customers basically using streaming on the Lakehouse architecture, um and these are some of the logos, some of the bigger logos here, but there are many customers. Again, this is kind of uh old data, there's even, you know, more customers more now using using streaming or Spark Streaming. So, let me talk about kind of um an overview of what you know, Spark Structured Streaming allows the user to do or some of benefits, key benefits here. So, one, one of the, I think, things that resonate with our customers and potential users the most for Spark Structured Streaming is it does have a unified batch and streaming API
For people that are kind of not familiar with Apache Spark, Apache Spark was this project that was originally a research project from a UC Berkeley, just not too far from here, right? A lot of the co-founders of Databricks were basically the original creators of Spark. And originally, it was kind of um a system designed to do more iterative processing over fixed amounts of data. But, it was then uh leveraged what then uh shifted to also be able to do stream processing. If you can think about it, you know, a little bit more, it's actually not that big of a shift. But, at the end of the day, the product of that was that you can have kind of some same or similar APIs to do both batch and stream processing. So, that was very kind of powerful for developers, especially for developers that are already familiar with Spark, to start writing basically streaming streaming pipelines or streaming queries, it became, you know, a lot easier kind of a a ramp up for them. So, another thing that uh Spark Streaming offers is fault tolerance and recovery. There's automatic checkpointing and failure recovery mechanisms that, you know, I can definitely elaborate more that allows for reliable operations
And we do uh provide exactly once guarantees um all the, you know, processing semantics for basically doing processing and structured streaming. And I would also say add on to that like compared to other peer systems, we have some other benefits in terms of recovery time. So, especially if you are interested in talking about like looking at instead of P99, which is the latency metric often used in stream processing, if you want to look even further at the Pmax, this is where recovery failure recovery really kind of factors in as well. Actually, there's a lot of benefits structure streaming has over uh peer streaming engines in terms of that. Performance and throughput, we see, you know, um a good amount of definitely a pretty sizable throughput that that we see through Spark Streaming handles, you know, 40 million events, 1.2 terabyte events per day, you know, for some of the more challenging workloads. And um you know, it handles it pretty well. So, definitely, we have seen use cases in which Spark Streaming for some of our bigger customers can handle definitely um a good amount of load. Flexibility in operations, um one of the things that Structured Streaming or Spark Streaming brings is you can not only express your processing logic in higher level um declarative languages like SQL or SQL-like languages, but also you're able to embed UDFs and arbitrary logic for both for both those very custom uh use cases, you can write arbitrary code, and the engine should be able to run, you know, basically run your your custom logic as well for you
Stateful processing, um we're able to support efficiently, you know, stateful aggregations, joins, and all the other, you know, plethora of stateful like deduplication as well natively in the engine as well. Those functionalities are built in. There's a whole list of them um that um you know, I can talk about later if you're interested or you can just look at look them up on the Databricks website or Apache Spark website. But, you know, stateful processing, there's quite a few, you know, built-in features that the user will not have to build. We also offer, for example, if there is something that is not a built-in kind of stateful processing operator or function, we have the tools and APIs for the user to build themselves. But, we try to have the kind of um the design in which the most common or the most useful of these stateful operators, we would already provide the implementation for you. But, again, if it's some very custom use case, there are APIs for you to use to express that custom logic. So, let me talk about some of the more, you know, straight the the the rise of some, you know, streaming applications with some new areas, right, that um I kind of also talked about in the streaming panel yesterday
And a lot of it ties in with we see use cases that require more lower latency, I would say. Before, for the longest time, latency was, I would say, comparatively, you know, higher. In the hours, in the tens of minutes, or even in the minutes now, right? But, now I think we have seen a rise in use cases that demand, you know, 99 percentile latencies to be perhaps in the hundreds of milliseconds or double-digit milliseconds, right? So, there has been definitely a shift um in expectations or requirements from users. And a lot of the use cases arise from kind of the more operational use cases, the, you know, where there might be some sort of person involved and kind of the pipeline of solving the, you know, the the problem. So, one, for example, is, like I mentioned before, proactive maintenance, like we see it in one of our customers in oil drilling, right? Um that that they want to be able to be alerted in real-time for imminent failures in the equipment, right? They want to have some sort of email or notification be sent out in real-time so that one of the engineers um or, you know, maintenance personnel could be aware of the situation and apply proactive like measures to mitigate the situation. Another thing, um you know, still down the lines of kind of proactive maintenance is, you know, elevator, you know, dispatches. First of all, monitoring kind of, you know, the health of the equipment, right? Being able to detect any sort of problems and also respond to kind of um different situations of how many people, you know, in the elevators and things like that. Again, these are kind of examples that I'm providing, but, you know, the general idea here is there are a lot of more use cases that are more operational in nature that require lower latency
Something different, I guess, is tracing microservices. Um you know, this comes down to one of the things I also talked about yesterday was observability, right? And a lot of microservices or, you know, other systems as well as how would you be able to monitor the health of basically not only software, but potentially even hardware in this instance as you want to be able to react to traces, how well things are doing, and um properly take kind of action to either, you know, correct something or just to kind of understand how well things are doing, right? So, tracing through microservices to maintain to make sure that, look, if you're also providing if the customer is providing some sort of end service to their customers to have actually metrics to know like, are we, you know, performing at our SLAs that we promised to our end customers, right? So, a lot of times for tracing through, you know, various kind of things and using kind of stream processing, you're able to, you know, monitor the health and performance of those systems in real-time. Yeah, like I mentioned, the kind of theme here is, you know, we want to have a lot of these use cases, we want to have latencies consist consistent with within like the sub sub mills sub second. Is kind of what we hear from a lot of customers that um we've talked to at Databricks. Um two, I think another kind of interesting thing is, you know, ease one of the other things that our customers have expressed is, can we basically have easier ways to express processing logic for complex use cases. Like I mentioned before, there are APIs in place for you to express that custom logic, but potentially another improvement is the user experience of that. It's often hard to use those APIs. It requires some level of expertise
Can we basically improve that user experience to for users to express some of the cost custom logic in a more kind of user-friendly and easier fashion so that an engineer does not have to get wrapped up for so long to do so. Um another is this comes back to um connectors and um it's always been an investment we've been making on my team is that um we are there are a lot of basically things we can connect to, right? And there's a dedicated team at Databricks to do so, to basically connect with all of the, you know, clouds, sources, and sinks, message buses, queues, and things like that. And that's also, you know, work we have been doing as well as like, you know, various other sources like Salesforce and things like that. Be able to kind of connect to various vendors and ecosystems for our people to ingest and process and also write outputs to. So, let me talk about kind of how, you know, structures how we have been working um and improving structured streaming to kind of satisfy the, you know, the the use cases and requirements that I mentioned before. So, one of the things that we announced um earlier kind of um earlier in the year is Project Lightspeed, which is a general kind of project that we have named to for all of the improvements um that we are making on structured streaming. And let me talk about kind of um what they are. One is performance that, you know, we want to make structured streaming and streaming faster, um low lower latency, more cost-efficient for end users
And the other side is also make it simpler for users to use. So, one of them, like I said, is predictable low latency. Our target reduction in tail latency is um you know, up to 2x. Enhanced functionality. So, advanced features like basically using um this is where we're working on new basically APIs for stateful processing to allow users to be able to use uh to write stateful queries in a more kind of um user-friendly fashion and allow users to express custom logic in a more useful user-friendly fashion. And also be able to for some types of for type some types of workloads like say custom custom windowing um for stateful processing, allowing users to be able to even express those um custom kind of functionality over custom windowing. We're working on that as well. Another big thing, you know, I mentioned that I think customers have mentioned to us is operations and troubleshooting
How do we kind of simplify deployment, operations, monitoring, and debugging? And, you know, we're, you know, actively working on um solutions and solutions to basically help with those as well. And Databricks has also invested in kind of other product lines to also help with basically these operations. I mentioned yesterday something like Delta Live Tables, which is kind of a serverless approach to stream processing in which users can kind of fire and forget. And, of course, connectors, it's always a big thing that um for customers, you have to for users, you have to be able to allow them to read and write to basically sources and so uh sources and destinations that they want to, right? So, we're always looking to add new connectors for our users to use to be able to, you know, connect with. So, let me talk about kind of predictable low latency and the improvements we've made there to um you know, Spark Streaming. So, let me talk let's, you know, first talk about kind of the motivation, right? Um like I mentioned before, there's an, you know, increased number of use cases that require real-time monitoring and operational alerting, you know, requires low latency. Um we define that to be under, you know, consistently under a second, but there are some use cases that could even be lower than that. Um operational pipelines typically have kind of the following characteristics: single-stage stateless pipelines, read from a message bus, writes to another message bus, right? This could be, you know, something basically Kafka some sort of processing to another Kafka topic or, you know, Pulsar topic to another Pulsar topic, whichever kind of message bus, you know, uh you're using
So, one of the things we wanted to do here is to improve kind of the bookkeeping or progress tracking in in structured streaming, which was on the critical path of processing. Thus, it was incurring, you know, kind of, um, a lot of latency. So, let me kind of go over how things were done before, right? Um, at the beginning, so Spark Streaming kind of takes a pessimistic approach in terms of processing that it wants to deliver correct exactly one semantics above everything. So, one of the things it does is it persists before even running the batch, it persists the it journals the kind of offset um, that it's going to process for this this batch beforehand. And then afterwards, it would also write to you it will write another entry into persistent storage marking the batch as done to basically, you know, offer very strong guarantees in processing of what we're going to process for this batch and what we've processed. This allows us to basically create a scenario that batches are deterministic and on reprocessing, the data will always be the same for this batch once the batch kind of is cut in a sense. Um, and the idea for Structured Streaming is that, you know, we can look at stream processing as in basically, uh, continuous stream of microbatches that you're just processing they you're you're you're processing the stream data within basically, you know, these, um, these units of microbatches and you're continuously processing basically newer microbatches as newer data comes in. So, when we look at these kind of progress tracking operations, um, they account for, you know, quite a sizeable percentage of kind of latency end-to-end latency
So, one of the things we did is just make the operations asynchronous, right? Um, to be able to basically, one is to make them not on the critical path, be done in the background, and as well as not as frequent so that they basically take no time on the critical path. Um, yeah. So, there are some implement implications on semantics and failure recovery. Like, well, if you don't checkpoint as often, you're going to take more time to basically, um, reprocess the data on failure, right? But, that is I think a tradeoff a user will, you know, in any system will have to make is, well, how pessimistic are you with failures, right? How do how often do you expect failures and how fast you want to recover? And we do have knobs for users basically, um, tune that to their specific need is how often do you think failures are going to happen? And if a failure happens, how long do you want to spend recovering or reprocessing data that you might otherwise have processed, right? So, those are knobs that and tradeoffs that users can make, but this functionality allows the this functionality this improvement allows the users to be able to kind of, you know, do that tradeoff between, you know, um, how much effort you're going to spend on the happy path versus how much effort you're going to spend on kind of the, you know, the failure path. Um, another kind of improvement we've done is, you know, log purging, especially for progress tracking is before we were purging all of this, you know, older entries of what we wrote to in the background, but now we've moved that asynchronously as well. So, we're just going to basically purge older entries in our progress tracking in the background, thus it does not account for, you know, basically any additional latency on the critical path, and that also, you know, improved latency a bit for us. So, let me kind of talk about the general, you know, at level of improvement we've seen. This is not exactly, uh, realistic benchmark
It was done on basically sort of dummy in-memory sources and sinks, but this gives you kind of the level of, um, improvement if we don't factor in external systems at all, right? Just within potentially the engine itself, what would be the level of improvement after we made these changes. And you can see like with kind of a batch duration where each how long it takes for each batch, which is indication of end-to-end processing latency, we've decreased that by like 10x by just basically making a lot of these things that are on the critical path asynchronous, all the book tracking things that or progress tracking things, uh, on instead of on the critical path, move it into the background and less frequent. So, it was definitely a huge transformative like, um, a change that that we've we've we've made there. And you can see like from this just this kind of benchmark, we were able to lower from an, you know, uh, batch duration each batch took about 337 milliseconds to around 30 31 milliseconds because, um, we basically moved all the the stuff that requires IO operations to persistent storage into the background. Another kind another improvement we've open a performance optimization or feature we've made is for this is more for stateful queries is, um, improved, you know, state checkpointing that we've all already we also moved, um, the checkpointing to be asynchronous as well. So, before, you know, every for stateful query for stateful pipelines or stateful queries, we were at the end of every kind of microbatch which is our kind of checkpoint barrier, um, we would basically checkpoint the state to external and persistent storage. But now, we have made the change to say, you know, we no longer just wait for that operation to happen before moving on to process new data, we've moved that into the background as well. Thus also basically decreasing, um, removing, um, uploading, the persisting the state checkpoint to to persistent storage off the critical path
So, that, um, you know, decreased our latency by 20 to 30% in stateful pipelines as well. So, um, there are some other more recent performance features, um, that we've also worked on that I will also worked on. Feel free to come and talk to me, you know, later. I don't think I have time. Like, one of them is, you know, pipelining. We're pipelining basically our execution which also is very transformative to our latency as as well. And we're also doing, you know, other things as well to improve latency. So, this definitely, um, an ongoing effort by, you know, Databricks and my team to improve the performance and latency of, um, Structured Streaming
But now, let me talk about some of the more like, you know, enhanced functionality that we worked on as well. That's another goal of my team is to, you know, work on improving the functionality of the engine. So, one is we want to make sure that Python is well supported of vast I would say a majority of our users are actually Python users. If you think about who uses these kind of systems is a lot of time data scientists or data engineers. Like, they're not going to be like they're likely not often going to be, you know, hardcore, you know, Java programmers, Scala programmers, and things like that. So, a lot of the use cases and a lot of users, like, what they're most familiar with is Python. So, we want to make sure that Python is well supported. So, you know, we've, um, there's already a lot of support for Python, but one of the goals is to make the Python support on parity with say our Java and Scala support which has been prioritized in the past because a lot of the users were more, you know, uh, more, you know, Java and Scala and technical people, but as we go forward, a lot of the users have changed to want to use Python
So, one of the things we've, you know, worked on is improving kind of expressing arbitrary stateful processing in Python and making sure that users using Python can express that logic just as, um, just as well as users using, um, Java or Scala. So, one of the improvements we're making is I kind of, you know, briefly touched up upon is the, uh, window functionality. Oh. Is, um, you know, what I don't have much time to go into this, but, you know, currently there's ways to do tumbling, sliding, and session windows, but a lot of time users want to define more custom logic on some sort of custom grouping of events. It's essentially what a window is. So, we're basically building APIs to allow users to basically trigger and go through the elements in a within a window, um, uh, in a custom fashion to allow users to express that kind of logic. Um, let me not go through this is kind of an example, but I'm kind of running out of time here. Another thing is state management
We're allowing users to basically, uh, read and also read the state in their queries from an external external, um, source and also be able to, uh, my ex also uh, evolve the schema of the state as well. As you would imagine, when data or date your data inputs change, the schema of your state could change. So, we're also working on basically allowing you to, uh, evolve the schema of your state as well. State rebalancing, um, something that we just allow is kind of a it makes auto scaling easier for users. We allow users to basically also allow the the engine would basically rebalance state automatically to new basically executors that, um, you're able that you new workers or executors you add to the cluster you know, this the system would a staple query would automatically be able to take advantage of those basically new executors so for auto scaling things like that we allow users. There's some also troubleshooting, you know, improvements we've made especially, you know, allowing users to export the metrics to a wide variety of, you know, other systems like data dog and things like that for them to monitor so we've you know, made those connectors available to customers and we're still going to work on more native solutions for this as well. Um we've added more connectors especially like Kinesis, Pub/Sub, things like that DynamoDB and we're going to keep adding, you know, connectors as well to structure streaming to keep supporting kind of ongoing use cases and future use cases for our customers. And there are many more things but unfortunately I don't think, you know, time allows and a lot of this stuff isn't open source so feel free to also, you know, just join the open source Spark community and you know, collaborate with us on this and we're intend to open source like most of the functionality we've, you know, um have been working on so we're definitely committed to, you know, open source Spark as well
Um that's it. That's it for me. I I think there's a lot more content, you know, I can talk more in detail especially some of the ongoing work we are doing ongoing work we're doing on, you know, the the engine but feel free to, you know, talk to me afterwards I can elaborate some more.