data.bythebay.io: Monal Daxini, Netflix Keystone - Streaming Data Pipeline @Scale in the Cloud
I hope you all had a good lunch and are trying to fight the food coma because I'm going to do a quick whirlwind tour through everything you've done in the pipeline over the past year I lead the stream processing effort in the real-time data infrastructure Matt Netflix many decisions business and both product are based on insights that together from the data we collected Netflix across all these different verticals including content product new talent infrastructure the whole gamut we started about a year ago to take over this old legacy system that we inherited and this legacy system started very simple all the need was to get all these data the events that were just starting to flow through our system and these an old project called Chuck while many of you may have heard of it or may not have heard of it and over time this evolved and evolved into a big hairball you have a lot of redundant data pads to get the same data over across many links and and sinks and this presented its own challenges every ridden Dan pathway had its own set of unique failures and it became very soon a nightmare to maintain and manage the system so this is something we inherited a year ago and we set out to replace this with a brand new system so that we could support at least ones processing semantics we could scale we could make it multi-tenant and unable stream processing as a service which I'll mention and talk about to it's the end of this talk and we also wanted to eclipse the dormant Chuck hua open source project that was heavily been used so our goal was to migrate this large pipeline transferring or one petabyte of event data every day in flight without any service destruction and the levy of your offered was not to lose more than point one percent of the data why we did this such as a huge challenge not only to switch teams sorry switched traffic in real time but also guarantee that data integrity stays put so early this year we launched the new pipeline looks like this we simplified a lot lot of aspects of it sorry the clicker stock listing so this was just two weeks after netflix announced its global launch in 2016 we launched to over 130 countries and it became a truly internet global TV provider and this means a lot of that's what we are slowing through our system sorry for the glitches but there we go let's just take a second to look at these numbers it's over 14,000 years worth of video that's being watched on netflix every day it's not minutes if not it's not ours it's actually years and so this leads to over 700 billion events that float through our system and we hit a trillion in December and right now we process over a trillion events every day accounting for all the fan outs that happen on towards the end we peek at 11 million events the average size ranges anywhere from a few hundred megabytes all the way to 10 megabytes and we process of our 1.3 petabytes a day so that I was 2.6 petabytes when we had two pipelines running an English fish over we still achieve four lines more than four lines right now availability of the service so if you want to know more about the evolution of the pipeline we have a blog post out there that you can take a look at so let's look at how the data actually flows through the system the event producers produce the data the hitter what's called the fronting kafka clusters they go through a sansa router and it ends up either all our events actually end up in s3 and optionally they end up in elasticsearch oak and zero Kafka depending on what the use cases the reason we split up our Kafka cluster is to better deal with availability and I'll go into a little more detail soon producing events all our event payload is immutable we don't touch the paler that the user gives us however we do have our custom metadata that we inject because we have our own binary wire protocol these fields are really useful in terms of doing analysis as well as trying to trace where they came from anarchism wire protocol is pretty efficient we currently support JSON ever us on the way we're going to support photograph and what this allows us to do is allows the stat additional metadata for traceability or other metrics they are going to want it to flow through it and also allows us to support different formats and make it both backwards and forwards compatible so that it's easy to upgrade our systems transparently without the users knowing about it and it's very efficient it's only 10 bytes overhead per message considering our message sizes are pretty large we wrapped our own Kafka producer on top of the existing producer so that we can make it more resilient and not have the application that's having a library actually fail in term whenever there's an outage right our number one priority is to serve videos flawlessly if that means we lose a few data points that we can analyze over you know that's fine but we don't want the experience to suffer so we do best stuff for delivery and react equals one and we do dynamic buffer size tuning on the producer so that we can minimize the loss and as you saw we don't lose much of the data at all and we have 49 availability in the amount of data we can put through the system so we have two major fronting two major classes of clusters one is the fronting cluster and the other one is the consumer Kafka clusters if you have a lot of data coming in and we want to make sure we don't lose this data if he allowed every consumer app to connect to fronting Catholic clusters then you'd have a situation where it could get overloaded very easily and that would result in a lot of data loss so we separated it what we do is wherever we have the fan out and you have a lot of consumer Catholic Luster's we're out those events to a different consumer graphical cluster and the only component that's allowed to talk to the fronting kafka cluster is our infrastructure so we can control and we know what scale is it's going to be hitting the kafka clusters with so that way we have kind of a good isolation between fronting Kafka and consumer Kafka clusters we start with 07 we are already moved to 09 and we moving to VPC we paid a lot of Pioneer tax and doing that in the cloud because it hadn't be done in the cloud of the scale we run it had we do a lot of open source contributions back either why are working with confluent or ourselves like the rack cover assignment was from our team and the other important aspect is we actually enable unclean leader elections because we care more about availability that does brings its own challenges of managing the cluster but gives us better to get availability so we run over 4000 Kafka broker nodes currently and the way via split them up is they have eight clusters in each region and we are spread across three regions so there are about twenty four clusters and for each cluster that's 24 clusters that we have 24 zookeeper clusters as well when we tried to launch the service last year we had a big meltdown on day one it was a weird combination of an edge case in zookeeper in Kafka and network issues is something that only happens in 17 thousand chance but it happened on day one will be launched it so that's where we moved and split to individual zookeeper cluster so that you have true isolation across clusters so a few quick tips try and stay under 10,000 partitions per cluster and under 200 nodes and you may want to leave at least forty percent disk space for their application traffic that it needs to happen on each broker so taking the whole available II in failure to the next level we built Kafka Kong so we actually every week route production traffic from one cluster to the backup cluster to make sure that when there is a real issue and a cluster is not responding you don't lose our data we just will switch over to a backup cluster that we build in real time we have the backup cluster size to only three nodes so it dynamically scales up the reason for three nodes is it lets us quickly create the topics while the other nodes are warming up and coming up and we can start taking traffic so this happens at least once a week it's part of the on call when we go on call the on-call person just does this and it's really simple you go into a tool we pick the cluster in the region we want to fail over and we just hit a button and it just happens so we have a prep failure button here what that lets us do is if there's a real scenario and the automation is not taking care of the drops then you could proactively go and prep a failover cluster and if there's a need we feel it over otherwise we try to recover the original cluster we also have an auditor process that runs out of band it runs in its own cluster and the benefit of this is it lets us monitor interesting metrics and report on it the reason of not having this in the broker process itself is if there's Network partitioning or any other issues that's not visible from outside this helps us make that happen and we can do things like you know performance monitoring and testing we have a metadata visualization tool that will open soon that lets us look at the health of kafka clusters and also search topics across all these clusters we just published a tech block as well he can go and read more about this if you're interested now it's the routing service the routing service is entrusting each other Sam sir because of the performance the resource utilization in the back burner show quoted had compared to spark one tattoo that I looked into last year our routing infrastructure is built on top of docker Samsa my sequel a little bit of go and see and we use goff car itself for the checkpointing cluster so we use the offsets that I read on the consumer side those get automatically stored in the checkpoint cluster and this is a built-in feature of samsar but we have to tweak a few things about it so on the high level we have on the control plane a router job manager that decides which docker container runs on which machine what it's going to which topic it's going to read from what partitions it's going to read from and what kind of filtering or projection or transformation is going to do so this is more a pipeline as a service that runs with the netflix and we set this end-to-end stream flow for our customers and once the data is in RDS about every job and the environment it needs to run in we have a small executor process that runs on actual physical hosts and it you know in a one-minute reconciliation look checks what's in the RDS database sees the docker containers that need to run on the host and runs those containers it only uses zookeeper one time only when it comes up or when and when that physical host actually goes down and that time it's used to assign an ID which is used to map to the jobs that it needs to run on so what we did is we run multiple Sam's our jobs for one Kafka topic that we are reading from and we run Sam's and a standalone mode Sam's as a job manager and a task manager and containers but we don't use the higher-level we directly control control each of the tasks and we run it in a standalone mode and we run every job that runs processes messages only for one sink for isolation if one sink is down if you have to process multiple messages and because Sam's I single threaded you'd be impacted together on one job that writes data to only one thing and you have one checkpointing topic / Kafka cluster now this is really important because if you have one checkpointing cluster too many checkpointing topics to the same source topic you're reading from now when you want to upgrade or do something else you worry about moving a state or if you change the partition assignment for the containers that they're running on then it makes it harder so this way if you have one checkpoint topic all the metadata about which partitions and how much progress you've made in the consumer side is in one place and you can dynamically scale the number of routing container so that's what we do if a topic gets more traffic less traffic we can shrink up and down the number of containers we use and this that's it makes it happen all our jobs once they start the configuration is not mutable it's immutable so we'd ever have to worry about it going wrong or some dynamic property change or changing things around and failing it's one less thing you have to worry about so this the executor I was talking about it also logs snapshots of a log that's running on the containers and in routine the uploads it to s3 and it also makes it available screaming so we have client tools available locally on our machine it's just a simple command line that you can log into any of these containers or you can stream the logs you can specify you want it logs from a certain period of time so it'll fetch from s3 if its historical it'll stream it will seamlessly stitch all these logs and give it to you so it gives us a very nice tooling to look at what's happening in our clusters so yes there's no messes running it's simple RDS and a reconciliation loop on on the node and this is great for us because let's moving pieces right so we made a bunch of changes to Sansa to make all this happen we use the thread job factory in production which is advised against in their docks but it works fine I implemented a fix for static partition assignments you can we can run it in standalone mode and we added both the regex and arrange partitioner specification there the other big problem running this at scale was Sam's allows you to specify a prefetch buffer with the counter in the problem it counts as you never know how much memory it's going to use because of the variable message size so put in a patch where you can cap the amount of memory it uses for prefetching and it dynamically manages memory so this these patches have been adopted in 20 10 and the samsa team bordered it over 20 10 I said performance really well it stays within ten to fifteen percent of what you specified because you know at the end of the dr running on the JVM and you don't know how much extra over it it's going to add but practically if you've seen it go no more than ten to fifteen percent and the overt of adding that was only point zero two percent so we've backported a few more configuration patches as well so that when we launched our containers we can specify an override environment variables to achieve the immutability of configs once the container starts and the sams our job starts the config don't change but we have we need a way to override it so we do it at the environment layer so we run fourteen thousand plus docker containers today to run the service and they run on over 1400 nodes across three regions and we don't do any cross region replication on our Kafka clusters it's all I lend one right now but we can route data from one cluster the other using the routing infrastructure where needed so give a talk last year at sams I meet up if you want all the gory details of what went in you can take a look at this link later will be posted so we measure metrics and all these data points that we have and we have a dashboard that's customer facing and what happens is anytime a new event is sent on a new topic that's been provisioned it automatically shows up here in the dashboard we don't have to manually enter it'll automatically show up all the routings will show up and the users can use it in a self-serve mode and we have a dell facing one as well so that we can look at our internal metrics that the users don't have to look at so we've scaled to you know 1,000,000,000 1,000,000,000,000 plus a day you know what so what we did is we expose the cost to our users and larger topics producers and what they found is one team found they were producing more than they actually need to they did some basic enhancements and they cut it down by six hundred percent we save costs overall really fast in the span of two weeks the costs attribution was a was a big deal as well and we do a lot of automation to free up our resources so that we can actually build infrastructure so everything you saw today is we build it and we run it we move to a more norms model than a DevOps model so we don't have any project managers we don't have any product managers we talk to the customers you're in infrastructure teams it's unique and we help drive the direction of Edwards it needs to go and the operations and n n and n to n so this doesn't mean we are we are overworked all the time we just make simple and better choices and we lean towards self healing systems so what are we doing in the future of your building systems so that we can improve the quality of the data whether it be in a schema registry discovery and and self-service tooling so stream processing as services as a big initiative that we are looking at we have different stream processing system that Netflix and we want to make it easy for the user to work on this so what we are doing is we are using you're going to use apache beam as our abstraction layer and you're going to have runners for different systems so that we can target and help the user right to one API and behind the scenes we can run it on different systems based on their strengths and weaknesses so this is how it's going to look like you know the user comes in creates a beamy API if you dr. eyes it for them you get it submitted they automatically get a dashboard and we get it up and running and they have an option of submitting a jar or a dsl and we take it from there so right now we're working on apache b'man fling is the next revision of our pipeline and the spaz offering and as we find more insight into this and scale as we'll share our findings so if you still need more data we have a bunch of links listed here there's a lot of material out there so feel free to read through or ping me or you have any questions and getting we're out of time yeah sorry we lost couple minutes in the presentation