Confluent Developer ft. Tim Berglund, Adi Polak & Viktor Gamov
Hi, we’re Tim Berglund, Adi Polak, and Viktor Gamov and we’re excited to bring you the Confluent Developer podcast (formerly “Streaming Audio.”) Our hand-crafted weekly episodes feature in-depth interviews with our community of software developers (actual human beings - not AI) talking about some of the most interesting challenges they’ve faced in their careers. We aim to explore the conditions that gave rise to each person’s technical hurdles, as well as how their experiences transformed their understanding and approach to building systems.
Whether you’re a seasoned open source data streaming engineer, or just someone who’s interested in learning more about Apache Kafka®, Apache Flink® and real-time data, we hope you’ll appreciate the stories, the discussion, and our effort to bring you a high-quality show worth your time.
Confluent Developer ft. Tim Berglund, Adi Polak & Viktor Gamov
Streaming Call of Duty at Activision with Apache Kafka ft. Yaroslav Tkachenko
Use Left/Right to seek, Home/End to jump to start or end. Hold shift to jump forward or backward.
Call of Duty: Modern Warfare is the most played Call of Duty multiplayer of this console generation with over $1 billion in sales and almost 300 million multiplayer matches. Behind the scenes, Yaroslav Tkachenko (Software Engineer and Architect, Activision) gets to be on the team behind it all, architecting, designing, and implementing their next-generation event streaming platform, including a large-scale, near-real-time streaming data pipeline using Kafka Streams and Kafka Connect.
Learn about how his team ingests huge amounts of data, what the backend of their massive distributed system looks like, and the automated services involved for collecting data from each pipeline.
EPISODE LINKS
- Building a Scalable and Extendable Data Pipeline for Call of Duty Games
- Deploying Kafka Connect Connectors
- Join the Confluent Community Slack
- Get 30% off Kafka Summit London registration with the code KSL20Audio
SEASON 2
Hosted by Tim Berglund, Adi Polak and Viktor Gamov
Produced and Edited by Noelle Gallagher, Peter Furia and Nurie Mohamed
Music by Coastal Kites
Artwork by Phil Vo
- 🎧 Subscribe to Confluent Developer wherever you listen to podcasts.
- ▶️ Subscribe on YouTube, and hit the 🔔 to catch new episodes.
- 👍 If you enjoyed this, please leave us a rating.
- 🎧 Confluent also has a podcast for tech leaders: "Life Is But A Stream" hosted by our friend, Joseph Morais.
How do you process all the telemetry data from everybody playing Call of Duty right now? We talked to Yaroslav Tikachenko about just that on today's episode of Streaming Audio, a podcast about Kafka, Confluent, and the cloud. Hello and welcome back to another episode of Streaming Audio. I am, as ever, your host, Tim Berglund. I'm joined in the studio today by Yaroslav Tikachenko. Yaroslav, welcome to the show. Thanks for inviting me. You got it. Now you work at Activision, and there's a lot to do at a game company. Always a lot to talk about, many, many diverse roles. But this being a podcast about Kafka, Confluent, and the cloud, I'm guessing uh that you are not a game developer proper, but have something to do with a data pipeline. Why don't you tell us about what you do there?
SPEAKER_00For sure. Uh so I'm a software architect in the Activision data department, which is an interesting mix, right? Because we have quite a few different data teams, and I'm happy to work with all the teams in the organization, but I mostly focused on the data pipeline uh development specifically. Um and it's all kinds of uh game telemetry from various game clients, various back-end services. We collect, we process, we store all the data from all those various games, primarily Call of Duty, and we have various data access tools for internal customers to access the data. And my role is um this interesting mix of being you know hands-on uh as well as driving some of the uh technical technical uh development and uh guiding the team in all things real time.
SPEAKER_01Awesome. And I think best kind of software architect, right, when you're hands-on, when you're still writing code. Uh if you're if you're really guiding a team or an organization as an architect and not occasionally building things, you can get a little unmoored from reality. And that's usually where architecture decisions become bad ones. So always good to see. Um, so give us an idea what this looks like. And I I sense we're gonna have to have this discussion kind of in historical terms, but probably everybody listening, whether they've played Call of Duty or not, understands what the experience of playing Call of Duty is like. It's a first-person shooter set in you know in various combat theaters across various uh uh time periods, um, you know, in its in its flavors. But what what what does the back end look like? I mean, what what's going on there?
SPEAKER_00So on modern online games, actually really, really complex beasts, right? It's just this really massive distributed system, uh, and there are so many layers. You know, you have the the actual piece of code that's running on the console, uh, as well as you know, all the online components, all the multiplayer components, uh, where we have a lot of back-end services for things like matchmaking, right? So, whenever you want to play a match, you need to find someone else who also wants to play a match at the same time, right? So uh all things around matchmaking, uh microtransactions, um, telemetry, statistics, leaderboards, like all those things. Um, in Activision, we have um a separate studio who handles all the online game development. And we have different game studios, you know, usually each year that work on the new version of Call of Duty. Uh, they all work with that studio on the on the new Call of Duty development. And so we have this layer of things like the game client, backend services, uh, some really interesting mix in between, which I am not supposed to talk about today. Um, but actually gathering data from all those things is a very tricky concept, you know, just because all the systems are so complex and they need to, you know, especially in the console environment, you have very, very strict memory requirements and CPU requirements. So it's very important for your telemetry code to not interfere with the game, you know, the actual game, but you also want to collect the data as accurate as possible. So it's it's this very interesting mix of things.
SPEAKER_01Sure. Because I guess you have like a frame rate contract in the console, right? There's you have a certain amount of time to get all of your computation done.
SPEAKER_00Exactly, exactly. And just because of that, uh we actually spend a lot of time in Activision building some of the in-house binary protocols. So binary protocols optimized for the game. Um, so for example, in in some of the versions, you know, we had uh fixed arrays. So the game could not produce dynamic arrays just because fixed arrays are simpler and uh you know cheaper to manage. So uh those are really tricky uh requirements to deal with.
SPEAKER_01Got it, got it. And so the the communication between the console and your back end is that custom binary protocol, and then you've got some proxy that takes that and I assume produces to a Kafka topic from there.
SPEAKER_00Exactly, exactly. So we have a REST proxy that we build in-house that writes data to Kafka, well, to set of Kafka clusters, because we have quite a few things in between. Um yeah, and we collect data from our data centers as well. Um, so we we currently use that REST proxy for you know all kinds of public um endpoints, for example, game consoles. Right now we're still collecting the game directly from uh from backend services to Kafka, uh, but we're also thinking about you know introducing that you know REST layer in between, even in our own data center, just because it can simplify a lot of things.
SPEAKER_01Cool. What is the I I I guess I want to start at the console. Um, and this is all console games, right? Um I want to start at the console, and you mentioned a little bit about that runtime environment, just that there's a hard real-time requirement. You know, you have a frame rate that you need to meet. And you know, in the game engine, if you've you I know you know this, but listeners, if you've read ever read anything about this kind of thing, there are things that game engines can do to you know decide objects not to render and detail not to include and things like that to make that you know, because you you don't you don't get to delay the shipment of the frame, it's it's going and you know you have a certain amount of time. So you've got this really constrained runtime environment. You said you do some some um unusual things to optimize execution time footprint. How many different kinds of consoles and different games are you mostly Call of Duty? And what what's the what's the spectrum of consoles that you have to support in that front end?
SPEAKER_00Sure. Um so actually I've just recently checked. Uh the oldest game that we still support in our pipeline is nine years old. Um you can yeah, you you can you can uh mix math and yeah, it's it's it's very tricky. So we we support the current generation as well as the previous generation of consoles, uh plus PCs, uh currently uh mostly on uh BattleNet, but before it was Steam, um, as well as for Call of Duty mobile. Right now we have you know iPhones and Androids. Um yeah, so that's um that's that's the list. And in terms of games, you know, Call of Duty obviously is our main main focus, uh, but we also you know shipped uh smaller titles like uh Crash Team Racing was a really successful launch this uh this year. And um, you know, so some older games like Skylanders, too. Still still running, still um ingesting data and using it. Nice.
SPEAKER_01So you've got this REST proxy, and uh that produces to Kafka topics. Take us from there. I mean, I mean we have uh talked a little bit about there being event-driven microservices on the back end, not on air, but you know, you and I have worked out some some things like there are event-driven microservices back there. This sounds like a giant streaming ETL problem, but uh boy, just walk me through what happens when telemetry lands in a topic. What do you do with it and how do you get it there?
SPEAKER_00Yeah, so um actually our end goal, just just to make things simpler, right? So from from one side, we have all this data coming through. On the other side, our end customers, you know, we do have some customers that are really comfortable consuming data directly from Kafka, and you know, they they're happy to deal with Kafka consumers and you know manage that that thing uh themselves. Uh, but we also have quite a few customers and quite a few use cases that just really want you know a simple tabular access using something like Presto, Spark, Hive, you know, uh very very classic big data tools, uh, as well as some of the NoSQL access, right? So we have some of the API-driven endpoints. So we have this uh set of interfaces we need to support, which is tabular, uh streaming, Kafka, and the kind of more classic uh REST API uh database um driven approach, right? And so we need to first process all this data that we have, and then we need to route it properly to all those different syncs. And we also need to manage schema evolution because schema evolution is very tricky in our case. And so uh right from that REST proxy, we we start thinking about what is the you know optimal set of topics, like how to even name a topic really, and why, right? Uh, because you need to understand how to scale that topic, how many partitions you need, how to, you know, how to route the data after that. And so we have um a service, internal service called refinery that handles you know like the initial data aggregation and routing. Um we we try to capture as as much information as we can in the initial set of topics, right? So where the data lands, but then that service, refiner service, the very initial kind of part of the ingestion, uh, we then uh rename not rename, but we kind of merge and and and split some of the topics based on some of the logical conventions we have. For example, in the most recent titles and in the newer pipeline that we have, you know, we try to really end up with a single Kafka topic for the whole game, right? Which might sound a bit crazy, but at the same time, we need to process the whole thing. So it kind of makes sense for us. And then if end customers need, they need a specific subset of the data, we have a service to split that topic even more based on the conditions they have, right? So the initial step is to route the data properly, uh ingest it properly. Uh then we transform it into the different formats, and we internally try to use Avro as much as we can. And uh we obviously, as I said, we have various um custom protocols, but we also support things like protocol buffers, JSON, uh some other proprietary formats, and we have a schema registry for all those things, right? So our data producers they upload schemas to the schema to the schema registry we provide, and then we use those schemas, we distralize all that into Avro on Kafka. Um, and uh from there we run uh a series e transformation steps, um, and then we expose the data to all those different interfaces.
SPEAKER_01The tabular and the API and the streaming. Exactly. Totally want to dig into that in a minute because that's there's some uh very timely topics in there that uh we've been talking about a lot recently. And uh recently I always like to make clear uh we're recording this in the middle of December 2019, so it'll probably ship uh early in 2020, but uh around the end of the year 2019 is kind of when we're thinking right now. Um now you talked about schema evolution and schema registry, and I didn't quite catch it. Is that confluence schema registry? Did you build a build one of your own? How do you handle all that?
SPEAKER_00So we we actually, yeah, we actually built a schema registry a while ago, um, many years ago, and there was no confluence schema registry available, I would say at the time.
SPEAKER_01It's the best reason not to use the confluence schema registry.
SPEAKER_00Yes. Uh I would say now, you know, why not? Why not use it? Um, except you know, over all those years, you know, we built support for you know some of the internal formats, some of the extra validation. Um, yeah, so it's just it it feels a bit more mature, just for our use cases, obviously, right? But it's all you know, it's all it's all can be done with the uh with the uh schema confluence schema registry. Cool too.
SPEAKER_01Um any special uh there's a there's a few things here. Um so I want to ask you about topic creation in a minute, um, but any special lessons learned about schema evolution, like what has worked well, what has not? Um I I wouldn't actually have guessed that that would be uh a daunting problem since the since game software doesn't change all that much. But I guess you have the goal is to get all of this into one topic. So no, I take that back. That's a terribly daunting problem. You have nine years of schema in there. So yeah, that's awful. Oh my goodness, seriously. Uh yeah. So how's that?
SPEAKER_00What have you learned there? And so the the really, really big problem here is we we're not only ingesting the production data, right? So the the data that's produced by the live game, uh, we also provide a service for all the Activision developers to send development data to us. So they end up with the you know streaming access or tabular access for the development data over the course of a year, for example, or a couple years when they develop the game. And obviously, when you develop something, you trade quickly, you end up just creating new schemas almost every day, right? Um, so the really important lesson that we learned is you you just absolutely have to automate it.
unknownRight?
SPEAKER_00So you need to come up with a set of standards, you need to come up with a set of tools, set of services that just handle it for you, right? And uh obviously then the question is how do you do that properly? Uh, but something you need to just realize immediately, you know, it's not going to be able to just um um resolve itself, uh, or you can't really um leave any kind of room where you need to do something manual, right? So it's all should be automated, uh, and you should provide a very strict guidelines. So, for example, um union types are very tricky to handle in a tabular uh way, right? So we don't support union types except like really, really simple ones.
SPEAKER_01Okay, well let me well come back to union types because it sounded like you described a topic with with fairly radical use of unions, but you don't support unions in the the events that get produced from the console.
SPEAKER_00They can't contain unions. We fine with simple ones, but we've seen use cases where people just try to put all kinds of stuff in this in this, you know, sure in a single union.
SPEAKER_01And it's I have 15 types and I don't really want to think about typing, and so let me make a union out of all of them. Okay. Now, um let me let me get to the union thing because it it if you've got uh you eventually want to get all of this into one topic, uh that's gonna be lots and lots of types, right? I mean, isn't that effectively don't you model that as a union?
SPEAKER_00So the topic itself, you know, we don't we don't really care that much about how many different types of messages we have in the topic because each single record in our topic has a schema, right? So any single message can have absolutely different schema uh as long as it has a schema ID and the schema is available in the registry, right? So we really use topic as just a way to transport things, but in the end, you know, each message can have absolutely different schema if if they really want to. Got it.
SPEAKER_01And uh I guess uh to set us up for the question I'm about to ask next, is there a language that uh most of your services are written in? You're consuming services?
SPEAKER_00Um so it's mostly Java Python, I would say. Uh we do support, we have some C sharp consumers as well. Um yeah, but it's mostly Java Python.
SPEAKER_01So as you're consuming from that topic with all of the types in it, you either reflectively or dynamically create the types that you need to instantiate objects of, I guess is how that would work. You have you have some wrapper around the consumer that does that for your schema registry.
SPEAKER_00Exactly, exactly. And we trying to push everyone right now internally to just use Avro, right? Because before we didn't provide a real-time streaming access to the data. Uh, and you know, each team that wanted to consume this real-time telemetry, they actually had to ingest parts and and you know, run some lightweight transformations on the data themselves. Uh because, well, we just didn't provide that real-time, um, real-time stream. Now, with the stream, we you know, no matter what kind of format we we ingest, can be protobuf, can be JSON, can be proprietary formats, we try to always expose data as Avro. So all the end consumers they only need to deal with Avro, and they all need, you know, they all know how to how to get schemas from the schema registry. So it's it's a bit simpler.
SPEAKER_01Yeah, okay, good. Um and all right, so we've got the the goal, there's there's various topics that things come into, and you have this goal of getting it into one uh kind of one topic to rule all telemetry. When you do create topics, though, you mentioned a little bit of a governance procedure around that. And this is a question I hear users of Kafka at scale ask sometimes, uh, you know, organizations like you. So if you would, um I have two questions. One, about how many developers uh might write code that functions as a Kafka producer or consumer, uh regardless of the level of the API, you know, streams or whatever, you know, any anybody who's concerned with topics, right? How many developers are doing that and describe your process for creating a new topic?
SPEAKER_00Sure. So historically, we actually try to go full autotopic create, uh, right. So we would let our data producers to create topics according to certain conventions, right? So if if there is a new completely new uh backend service in Activision, uh, or maybe there is a new uh game ID, right? So uh we kind of partition things initially that way, um, they should be able to just start sending data to Kafka without any kind of extra permissions or any kind of coordination. They want to send us data, we should be able to process the data. And this is very, very core principle we've been trying to follow. Uh, no matter what data producer throws at you, try to ingest it and process it, right? So we've been very, very flexible in terms of how many topics, what kind of topics people want to create, and we just allowed that, allowed them to send the data. Uh, but after that step, internally we look at all those different things, all those different conventions, and now we're trying to introduce more standards and to introduce more uh you know like fat topics comparing to the previous set of slim topics, because actually the number of topics and number of partitions you need to manage, you know, start to matter a lot when you reach certain scale. So we're trying to now come in this uh in this problem where we have just a few small um a few number of really fat topics to manage uh in comparison to a very large number of topics before, where each topic can produce just you know a very small volume of data. So it's it's kind of you know trying to inverse that right now.
SPEAKER_01Um and how many, I don't know if you're able to disclose, how many people would be thinking about this? Would be potential sources of new topics?
SPEAKER_00Oh, so uh we try to abstract all this from the you know the actual data producers. So um I would say we have what five, seven teams um working on the different back-end services. Um, and then couple studios usually actively developing and uh developing the game, but um again, they they just usually uh interface with the REST API or some of the SDKs we provide. So they they shouldn't really care about you know what kind of topics they create.
SPEAKER_01Right. They are in that case consumers of other interfaces and not consumers of Kafka directly. Exactly. Yeah. I I ask because there's there's I think a sliding scale of governance, you know, once you get to say a thousand people using your Kafka clusters, and you know, somebody out there wants to create a topic, you just don't have you know, as as the the person responsible for the service or administrating the cluster or whatever the right word is, depending on how you set things up organizationally. Um you don't really have a lot of common interests with random person from that thousand count pool of developers, right? You work for the same company, but you're you're as good as not, you know. So there tends to be a lot of governance at that scale.
SPEAKER_00Um and but at the same time, you know, in the last two years, I've been hearing self-serve everywhere. And uh I love the model of self-serve UIs, where you can expose some kind of interface to your end customers, and maybe you can also run some kind of capacity planning or some some some cost numbers right there. So when people decide I need this you know, huge topic with this, you know, with this number. Partitions with this uh messages uh per second, they can see immediately how much it's actually you know gonna be costing for them and uh if we have enough capacity at all, right? Right.
SPEAKER_01And you know, I I want to hear the word self-serve more. I really do. Everybody do more self-serve if you can. Uh but it it seems like at scale those those bolts always get tightened down a little bit more. But you're at uh still a pretty workable scale in terms of the number of human beings who might actually want to do this. So it sounds like you've got something that's still relatively free and rewards experimentation and rewards innovation. Um, and you've got a small enough group of people that you know there's a sense of common interest and and community around this thing that tends to keep those decisions uh being good decisions.
SPEAKER_00True, because uh I I notice you know some companies they they now offer Kafka as a service, right? So they just provide Kafka as as a as a piece of tech that people can use in in you know in really in in different scenarios as different, you know. Um you can build different solutions on top of Kafka. And and we try to provide a fixed number of solutions, a fixed number of use cases for Kafka, but we allow people to send uh different kinds of data, right? So uh we really don't try to um position ourselves as Kafka as a service team. We are a data pipeline, uh a data platform. So we we're trying to provide different levels of interfaces, right? REST proxies and SDs.
SPEAKER_01I that is, and Yaroslav, that is absolutely the right framing. You're not Kafka as a service. I mean Confluent Cloud is Kafka as a service. Um you are you are um a data layer for online games at Acavision. That's right. Um and speaking of which, I I don't uh think I asked this. Do you guys run on-prem? Oh, so I'm sure how much I can do that. Oh, totally cry. That's fair point. Yeah, yeah. I just want to know everything. So just give me all the details, private keys, you know, whatever.
SPEAKER_00Uh I I definitely can say a few words. Um, the game itself, uh, all the back-end services, obviously. You know, we have a few data centers, uh, but we also run a few things in cloud. Uh all the data stuff though, like all this data pipelines we run on ALES, um, they've been really good partners to us. Excellent. Um, it's all it's all Amazon's called.
unknownYeah.
SPEAKER_01Uh okay, uh these other things. It's kind of funny. Like you said, I don't know, three minutes of stuff at the beginning, and I am still unpacking the footnotes in that three-minute statement. Uh you mentioned three kinds of access to the data, uh, tabular, API, and streaming. And you said some other kind of compelling things in there. So could you possibly give examples of data? And again, I I will I will always ask questions very likely beyond what you can disclose. So anytime I do that, just don't obviously just tell me that you can't. But um to make those clear to the listener, I have a vague sense of what you might do with those three things. But if you can give examples of tabular data, API data, and streaming data, I want to dig into those because there are some interesting emerging architectural ideas lurking in there that I want to see if we can draw out.
SPEAKER_00For sure. Um, so in terms of streaming data, I have my favorite example that I always give to people. A couple years ago, we launched Call of Duty World War II, and we also launched Alexa app for Call of Duty. And so we actually have a team in Activision that built this really interesting machine learning um system where they would in real time um ingest the data that we already processed, uh, and then they would build certain uh models. And if you have Alexa, after you finish playing the game, you could ask Alexa, hey, how did I do? And Alexa would tell you, well, um, you know, next time use this weapon and use this loadout and you know suggest a few other things, and also say, by the way, but you're still better than you know most of your friends, and this and this is really cool. Like this is really uh cool real-time use case. Yeah.
SPEAKER_01And but that's that Alexa skill is making an API, there's a synchronous API call behind that skill, right?
SPEAKER_00So what they really wanted to achieve, they wanted to shoot this real-time machine learning, right? So they wanted to ingest data and also build certain models in the background. And the team also manages some of the you know kind of personalization services. So if you, you know, um uh if you uh reached, I don't know, 20 kills in the game, they want to be able to know that information immediately so they can show it on the I don't know, um uh mobile app, web app, as well as use it in that uh personalization engine. So they they kind of trying to decouple uh calling all the data services versus the you know the data models they have themselves. Uh and it makes sense because they want to optimize things certainly.
SPEAKER_01So let me dig into that a little bit more. You still and actually let me let me tell you about the distinction that I'm thinking about. So uh so you don't have to guess. But there are uh there are some kinds of consumers of streaming data that um want events pushed to them, right? Like I'm I'm doing some computing in aggregate, I'm running some machine learning model, and I want each new event pushed into me, and I will do my computation and my state will change according to that computation, and and we'll go from there. But then there are other consumers of streaming data that actually want to go and ask, and it might be like the state of that model or the the table that results from an aggregation. Somebody might want to, as it were, synchronously say, hey, you know, for player owned you 372, um just made that up, maybe that actually is a name. Apologies if it is. Um, you know, what what's the current score, what's the current kill count? You know, that would be an aggregation, a windowed aggregation um over some like session window kind of thing. So there's there's queries where you go and ask synchronously, and then there's um stream processors that need messages asynchronously pushed into them. And you listed tabular API and streaming. And I just want to try and figure out like, okay, so streaming would be these ML models, right? That would be a streaming consumer.
SPEAKER_00Yeah, exactly. And any kind of tooling that really needs the you know the access for data in real time. For example, some of the game studios, they would kind of reconsume the data again that they send to us, and they want this really low latency dashboards in the studio so they can you know make make changes. And this is you know uh a really good use case for something like metrics, right? So they would send uh a bunch of metrics from the consoles, things like error rates, things like uh you know, crash dumps, and they want to see the data immediately as soon as possible in you know all kinds of dashboards and visualization.
SPEAKER_01And I'm I'm envisioning some sort of network control center or network operations center, sort of big dark room with lots of screens. Okay. Exactly. Yeah. Uh that word picture is amazing, by the way. That sounds like a very fun place to spend time. Um I'm sure highly secure, but it would be cool to be on the inside. Um you mentioned some tabular or NoSQL things. Can you can you talk about what goes on there? Because that's also sort of interesting to take streaming data and and put it in a big database and then and then have some process that queries it. So what's the use case there?
SPEAKER_00Yeah, so so tabular is actually our you know one of the core uh products that we provide. So we spend a lot of time on making sure the thing is optimized. And um the architecture is very simple, uh, right. So we have uh Apache Hive Metastore uh with a bunch of tables on S3. So our goal is to uh archive the data on S3 in the right format and also modify the schema in Hive to reflect you know new fields, new tables, uh, all the you know standard schema evolution uh tricks. And uh the interesting thing is uh this year we actually switched to the new version of the pipeline where we started using Kafka Connect for the S3 integration, and we invested a lot in trying to make that you know really seamless experience for all the developers and for the end users. And we also invested a lot of time in automation on the you know on the schema evolution side on the uh on this hive schema side. So we have services that will automatically understand you know uh if um due to a schema update, you need to append uh a few new columns to the table in Hive. So they just you know deal with the schema evolution again in in this very automated way. Uh the same thing goes about tables, uh new tables, new events, uh new partitions, right? So we we're trying to to use connect uh with S3 to uh partition and transform and you know uh upload the data to S3. And we also have set of services for the uh Hive schema evolution.
SPEAKER_01Awesome. Awesome. So uh for some of those tabular views, you end up with data in S3 and then and then uh services in front of that that uh allow you to query.
SPEAKER_00Yeah, we use uh Presto Spark, and um again, this this this use case is very important because all the reporting is you know based on that. So if Activision Leadership wants to know how much money we made last week, this is where you know ultimately you would uh you would go looking.
SPEAKER_01Oh, nice. Okay, that's quite timely. Tell us, and you've been there two and a half years, so what what is the before and after? And and this is a typical concern of the software architect. You've kind of hinted at some of this, but if you could just step back and say, you know, this is what existed uh some number of years ago, and this is what exists now, and this is the future state you're moving toward. Um, what's the the kind of timeline of all this?
SPEAKER_00For sure. Um, so the pipeline itself is very old. You know, we survived quite a few iterations. And uh the previous iteration we still, you know, we still use, we still run uh quite a few services. Um it was a very classic ETL, you know, Hadoop stack ETL, where we have data on S3 and we use MapReduce and Hive and Spark and all those classic big data tools to actually run ETL. Uh but for that pipeline, uh, the interesting part was the ingestion of data, because the data originally landed in Kafka, but Kafka was managed by a different team. Um, and then when this team tried to mirror the data and process it and run some transformations, uh, the team actually faced quite a few limitations due to the way they run Kafka. So I cannot really uh tell too much about how exactly uh it was done, but the summary was um you just could not add more brokers or partitions after a certain, you know, certain certain level you reached due to some licensing and and some other uh interesting restrictions. And so the team got very creative, and you know, you have Kafka, you still need to scale. So they ended up using SQS. This is the Amazon's um uh queue service, right? So we still ended up using Kafka as a message bus to pass you know certain metadata messages around, but for the actual kind of you know work um work distribution, um, we would end up using SQS. And so it was a very tricky mix of Kafka and SQS, and then you're going into this Hadoop world, right? So the pipeline was quite complex. And uh the tricky part again was handling the scale, uh, because for production scale, you know, this makes sense, and people are happy to wait for you know four to six to even 24 hours while all this MapReduce and hive tasks are running. But you also have this development data use case, and people really want to see development data within seconds or at least minutes. And so we had to build like a smaller version of that ETL just for development data scale. Uh, but that meant a different set of tables, different set of schemas. So people got really unhappy with the way you know this distinction uh happened, right? So when you start developing a game, you would end up querying and accessing uh one set of tables. When you launched the new game, you would switch to you you would have to switch to a new set of tables. Uh and that was again really frustrating.
SPEAKER_01And error prone, sounds like, for transitioning from deployment or from uh development to deployment. Yeah, yeah, definitely.
SPEAKER_00And you know, this year we really invested in a lot of infrastructure, and we just decided, you know, we want to go 100% Kafka uh using Kafka streams, Kafka Connect, uh using as much Kafka as we can. Um and yeah, so right now we we don't have that layer of classic Hadoop ETL anymore, just because we ingest and we transform all the data using Kafka streams, and we use Kafka Connect for certain integrations like S3, like Cassandra, like Elasticsearch. Right? So all that is handled in Connect, and we are trying to uh use streams and streaming ETL as as much as we can. And this actually dramatically reduced latency. Uh instead of waiting for four to six hours as an end customer, now you wait maybe up to five, ten minutes in a production scale, which was really nice change. And apparently the pipeline is also way cheaper, but that's mostly due to the laziness of us because we didn't probably spend enough time trying to optimize the the Hadoop um ecosystem, right? But it's still really, really uh, you know, big difference in in the cost as well.
SPEAKER_01Cool. And uh on that last point about your uh your alleged sloth, uh, you know, fair enough. We don't always optimize everything, and then I think it's it's usually good judgment not to optimize too soon. But I mean it is a win, it's still a win. I mean, uh you were, as you say, too lazy to optimize Hadoop. I don't know if you've uh suddenly discovered uh this zest for optimization and you've squeezed every cycle out of your Kafka cluster. If you can take basically the same operational approach and get a cheaper footprint, you know, cheaper computational and storage footprint, then that's a win, right? That's a good thing. Exactly. Yeah. Streams. Talk to me about streams. That's uh fantastic. What do you do with streams? So we have this classic. Well, I guess you said you do everything with streams, so you've already answered that question. But what I mean is tell me about some uh you know, some transformations and just your general approach.
SPEAKER_00Sure. Uh as a data platform team, we really try to focus on the you know transformation side of things. We're not trying a lot of streaming analytics or machine learning, like we have different teams. Uh, but uh ourselves, we invest in uh you know Kafka stream services that ingest data, parse it, you know, they they can access the schema registry. We're trying to transform everything into Avro, and then we run some transformations, for example, for the tabular representation. You know, when you have this deeply nested data, like you can have up to seven seven levels of various objects and arrays in certain payloads. Uh again, it's really tricky to display that data in tabular ways. So we have services that can flatten it, that can create certain child tables with the uh you know, certain columns that it can use to join things. And so we try to use a single service as a big transformation step and then just have different stages of the pipeline represented uh as Kafka topics. So if uh a single, you know, if if we have a customer who's interested in a certain intermediate stage, and we we certainly have those use cases, they can just consume those topics, right? But if they care about the end result, they can either consume the final, you know, tabular topic or the actual set of tables.
SPEAKER_01Awesome. Awesome. All right. Other things that come up in your world that are of interest to I mean, there's just a there's a few things that seem like they would they would happen in game telemetry uh that people always ask me about, like um personal data, you know, PII, things like where where does that where does that enter your world and what sort of processes or uh frameworks have you put in place to help?
SPEAKER_00Sure. So we try to always follow certain uh certain standards, certain guides that we basically created for ourselves when we handle PAI. And we actually have uh you know an actual service that can handle the PAI kind of uh cleanup, uh how we can call it, or um you know, extracting the relevant PAI and then anonymizing data for all the downstream customers. So we we try to make sure we always capture all the relevant PAI in case we need it, and this this you know set of tables is you know audited, it's plugged down, um, and the data is always encrypted and you know all the all the proper words I'm gonna say. Um but for the end users, you know, usually, usually as a you know, as an analyst or data scientist, you don't care about PA that much. And so we're trying to uh you know always anonymize data. Uh and again, it's just due to a certain standards we uh we established, like if you really need to send a PAI to the you know the data lake, uh you need to also provide this you know set of uh set of fields. For example, you know, with the latest uh Call of Duty, we now have crossplay, where um people from PlayStation can play with people on Xbox, right? And if you have, let's say, a username, you really need to understand if it's a username on PlayStation or Xbox, because you know the same username can can exist in both platforms. And you know, by requesting uh really uh you know simple rules, but um nevertheless, rules uh from our data producers, we can gather all the metadata that's necessary for for us to anonymize the data properly and then also have a way to access it if needed.
SPEAKER_01Now, in the transition from the old Hadoop ETL world to the current uh streaming world, um uh uh across uh large relatively large number of years of game versions and console versions and so forth, um my guess is that you had uh certain internal customers for data uh that you computed in the old way that you migrated to the new way, right? It wasn't just all standing up new greenfield data sources, you had to migrate things. Definitely. How did you go about that? Uh I mean it's this is I guess this is kind of a special form of the how do I refactor the monolith to services question. It's kind of like that. Um, but how did you do that? Like you know, you did you keep the external contract and the customer knew nothing? Did you change that API? Uh walk me through one of those migrations.
SPEAKER_00So uh the biggest the biggest use case to discuss here really is tabular access, because you know, if it if it's streaming access, then you just end up with a different set of schemas, right? Right, said suddenly. And if you have that schema integration um working, then it will just handle it. So this is this is not very, very hard to manage. The same about API, right? Uh if your RESTful API provides a very generic uh message representation, for example, it returns a certain envelope, and then inside that envelope you have you know any kind of arbitrary set of bytes, and then again the schema maybe, or it's set of parsed bytes in in in JSON or some other format. Like this is all very manageable, I would say. Uh the tabular access though is is tricky because people run a lot of uh ETL on top of that. People use a lot of dashboards on top of that, and people in general were used to the idea of you know, if I need to calculate the number of daily active users, I just need to run this piece of SQL and that's it, right? So it's very easy with tabular access, you suddenly have way more users because it's just simple SQL, and people suddenly start to you know to be attached to uh if you very important tables. So any kind of attempt to change that, you know, is uh is very tricky because people just just asking, why would you need to change that? I'm you know I'm I'm so happy with what I have. So the strategy is very simple in this case. We provided a set of you know new schemas in in Hive MetaStore with new tables, and we said, you know what, uh you can still keep using you know the old way for now. We allow you that. But if you switch to the new set of tables, uh the rate uh of updates on on the data is you know 10 minutes versus six hours. There you go. And also we handle the duplication better due to all those extra things, and also we handle small files problem better uh for the development data, and they're like, hmm, let me try that. And uh, you know, we we've got a few early customers that were really surprised by the performance, by latency, and and and yeah, it's it's it's not a not a quick exercise, definitely. It's it's a long process. But if you give your customers something that's you know significantly better and um you you push them a little bit to the right direction, um, they will figure it out.
SPEAKER_01Yes. Uh I can't underscore that enough. Everyone listening, um, that's how you do a migration. You produce a better service and then let your customers volunteer to upgrade to that service because they're getting value that they couldn't get in the old version. Uh you don't Hector them, uh, you don't you don't uh use a club, uh you give them something that works better. And um how long, how long did it take on average, if you can recall uh a few transitions? It's still happening as usual. So you know, yeah, people have their own struggles, and maybe the maybe the old version's good enough for now, but um that's uh that still strikes me as a brilliant way to do it. Thanks. My guest today has been Yaroslav Thachenko. Yaroslav, thanks for being a part of Streaming Audio. Thank you. And there you have it. Hey, it's Kafka Summit time again, and you get another discount code for listening all the way to the end. Kafka Summit London is coming up on April 27th and 28th, 2020, and you can get 30% off your registration if you go to Kafka-summit.org and use the discount code KSL20Audio during checkout. Just enter KSL20 audio while registering at Kafka-summit.org, and that 30% off is all yours. I would love to see you there. And anyway, I hope this podcast was helpful to you. If you want to discuss the podcast or ask a question, you can reach out to me at at TLberglund on Twitter. That's at T L B-E-R-G-L-U-N-D, or you can leave a comment on a YouTube video or reach out to us in Community Slack. There's a Slack sign-up link in the show notes if you want to join that group. And while you're at it, please subscribe to our YouTube channel and to this podcast wherever fine podcasts are sold. If you subscribe through iTunes, be sure to leave us a review there. That helps other people discover the podcast, which we think is a good thing. Thanks for your support, and we'll see you next time.