Speaker Application
SPEAKER APPLICATION
Sponsor
SPONSOR
Past Events
PAST EVENTS
2026 SF
2025 OAK
2024 AUS
2023 AUS
2022 AUS
2020 SF
All talks
A
L
L
T
A
L
K
S
Building a Music Analytics Pipeline at Pandora
TBA
Speakers
Brian Femiano
, Senior Data Engineer, Pandora
View transcript
00:00
My name's Brian. Thank you, Rick, for introducing me. I've been doing data engineering for about
00:15
10 years now, different industries. I contributed a little bit to Hive back in the day, back
00:21
when it was still MR, and a project called Giraffe, which many of you probably haven't
00:27
even heard of. But for now, I work on the artist, creators, and tools team at Pandora.
00:33
So we work with social media and the Pandora data to try to create useful insights for
00:38
artists, labels, managers, pretty much everybody but the listeners.
00:47
So every day, we send the music labels a data set containing song play events. So what do
00:52
I mean by a song play event? Every time somebody listens to a song on our platform, we report
00:56
on that. So we call that a spin internally. Kind of makes sense when you think about it.
01:04
And there's actually people from the labels here, the data engineering teams that receive
01:07
that data. So I was not expecting that, but the pressure's on, I guess.
01:14
So why do we have to do this? I mean, the short answer is we have to, right? But that
01:20
doesn't help. Why do we have to? So a bit of a history. Pandora has got a reputation
01:27
for a very personalized music streaming service. People like our recommenders. But I think
01:32
the industry has proven out in the last couple of years, especially, that people really want
01:37
on demand, too. Now that that's possible, and a lot of the legal issues have been worked
01:41
out, and those battles have been fought, people like that. I might even wager to guess many
01:45
people in this room probably have a Spotify account, or Apple Music, or something other
01:49
than Pandora, even if you're a Pandora listener. So they found, Pandora found that most of
01:54
our users were complimenting Pandora with something else, and still do. Like, they hear
01:58
a song, they really like it, and they can't just listen to it again, so, even if they're
02:03
paying for it. So they go to YouTube, and they would look it up, and we just didn't
02:07
really want that. So now we actually have an on demand streaming tier. There's different
02:11
ways to access it. You can be a paying customer and get access to that. I won't go through
02:15
all that today, but in terms of actually, so now that we support that, we have to report,
02:21
labels want reporting on every single instance of a song spinning, for which they have rights
02:28
to. So we have to send those data sets to them every day. And it's not even a report,
02:32
per se, like this song spun X number of times, they actually want a structured event for
02:37
each individual spin event that happens. So these get to be big, and we have to, they're
02:42
external data sets that we're effectively delivering over S3 every day. And it's a new
02:49
requirement. Pandora didn't always have to do this. We had kind of an adversarial relationship
02:52
with the labels for years, and that all had to change very rapidly when we wanted to do
02:57
on demand, because they have the rights to that music, so. So what do these events look
03:02
like? This is a very stripped down version. There's about 20, 22 different structured
03:06
fields, but they go something like this. This kind of gets to the core of it. Like, what
03:11
time did it play down to the second, right? What's the name of the track? Some others
03:15
you could predict, right? Maybe the artist's name, I didn't include here. How long did
03:18
the track play when it spun? That's important. And then I, one of the more interesting ones
03:24
I threw on here is the play source. So that gets, that's a very complex field internally
03:29
that we manage. How did the user actually launch the song on the app? There's dozens
03:33
of ways this can actually happen. Maybe it was just a traditional station spin from the
03:37
recommenders. Maybe it was they have, like, a premium playlist, and the playlist ended,
03:43
and so our recommenders kick in there, too, and start, it seeds everything off your playlist
03:48
and starts spitting out new songs. It just plays them. Or maybe they were, like, a premium
03:53
access thing where they're not paying for it, but they wanted, they were willing to
03:55
listen to a 30-second video ad, and then we give them access to the song they search for.
04:00
That ended up being wildly popular when we launched that. And so where does this data
04:09
come from, big picture? So there's lots of, Pandora's a huge hive shop. We run all sorts
04:16
of different Hortonworks clusters, completely isolated, but not totally firewalled off,
04:21
but the data that we care about comes from one of our primary analytics clusters. So
04:26
there's a huge trough of partitioned spins data, all in parquet form, that we generate
04:33
off of their raw data, and then we have metadata that we build up, too. So I only included
04:39
two different tables here to show the metadata, but we consider all single station ads, every
04:45
single station that a user has as metadata across all our users. It's a big data set,
04:49
but we call it metadata because relative to the spins data and the fact-level data, it's
04:54
not that big. Listener data, and there's others like artists, track roles that artists
04:59
had in the track, composers, album, all of that stuff. And listeners is an important
05:06
dimension, too. So that comes in every day from external sync processes, which my team
05:11
doesn't have to think about. And we, the whole pipeline that I'll show is when we detect
05:18
that our transformation succeeded on the analytics cluster for a given day in all this, we just
05:23
CP that back to a private cluster that we can control, so we can rely on it to do this
05:27
reporting that we need to do. And then also I threw on here, we do, Pandora does, we have
05:34
Kafka queues where this stuff all lives on, too. Our team has not gone through the motions
05:38
to actually transform or update our report to actually use that to do this reporting.
05:44
Just because we haven't gotten around to it yet, but we, they do throw spins and other
05:49
things on the, and other metadata on the Kafka queues we have. We just, our team doesn't
05:53
actually leverage that yet, but we plan to. It would cut down the latency. So the first
06:00
iteration of the entire pipeline, which I'll show in a second, it was pretty slow. I mean,
06:07
to get data from like all the spins that happened on a Tuesday over to the labels would take
06:11
to like Thursday. And we just wanted to do better than that. So we wanted to improve
06:16
the overall latency of how we send these files on S3. And to do that, we took, we looked
06:21
at two primary things. One is the obvious one. It's like the runtime of each individual
06:25
step in the workflow. How could we improve that? We're glaring issues. And then also
06:30
the bigger one, more indirect, is how do we recover from issues that happen along each
06:34
way? Like if a disk CP fails or something like that, we don't want to have to redo upstream
06:38
work. And so we set out to break out all those checkpoints. So as you can see, I kind
06:49
of structured this whole presentation as kind of like a war story form where I actually
06:53
cut out a lot of the like, well, this is what we used to do and we didn't like, and I'm
06:57
just going right to the solution. So I don't have a lot of like before and after runtimes.
07:01
There's plenty of time for questions. I'm sure some of that will come up, but I don't,
07:06
it feels like it doesn't really help to brag about those numbers either because it's not
07:08
that interesting to people outside the team. So what were issues with the previous pipeline?
07:15
Just to summarize, we'll give it one slide for that part. Well, there are upstream delays
07:20
on the raw data coming into us. We don't really control a lot of the sync processes or some
07:24
of those other analytics teams and what they do. And there's just a natural delay because
07:29
of how they're working with daily partitions. So some of that's delayed and we just can't
07:34
control that. But the data set generation, RN, once we have it in that private cluster
07:39
was very slow. And I won't get into a lot of reasons why, but it just was. And then
07:45
also part of our requirement structure to the labels is we send a single file per day.
07:52
And I've asked around like why that is. I just accept that it is for now. And so I know
07:58
it's a bit of like a big data anti-pattern, right?
07:59
A lot of us were used to working with directories of files
08:03
or buckets of files, where all the files represent
08:05
everything in the day.
08:06
And you get better parallelism that way.
08:08
But we send a single file per day
08:11
to each label for all the spins that
08:13
happen for tracks that they have rights to.
08:16
And then so part of the problem was the steps two and three
08:18
where we generated and we're making the single file
08:21
were tightly coupled.
08:22
So any kind of issues with either or
08:24
would scuttle the entire thing.
08:27
So the new pipeline follows just a series of linear steps.
08:31
Upstream on the analytics cluster,
08:33
we do some transformations for both the raw and metadata.
08:36
And I'll cover a little bit more of that in detail,
08:38
like what that means.
08:39
But we do that work directly on the Pandora analytics
08:42
clusters through with Hive on Tez.
08:45
So we use it that way.
08:48
Once that data is prepared and those partitions are active
08:52
and we validated them, we distribute them
08:54
back to our private cluster in batch.
08:57
And then once we detect that the atomic DCP has succeeded,
09:01
we can begin building.
09:02
We use Spark very heavily on our team
09:04
to generate this data set.
09:06
It's one of the uses for it.
09:07
It was one of our first uses of Spark on our team, actually.
09:11
This delivery was so important that we set out to.
09:14
Spark was actually one of the early test cases for it.
09:18
And then so once we have this data set,
09:20
and the data set itself has been validated,
09:21
we want to coalesce it into a single file for the labels.
09:24
And then we just have to push that up to S3.
09:31
So in order to automate all this,
09:33
we kind of codify each level, each step of the workflow
09:36
in that five-point series.
09:38
We use Python classes.
09:39
Our team still uses Luigi.
09:41
I know a lot of the mind share in the data engineering space
09:44
has kind of moved over to Airflow.
09:45
And other frameworks are more popular.
09:47
But we still really like Luigi on our team.
09:50
We feel like it's lightweight.
09:51
The visualizer works very well for us.
09:54
And I think the most important part is the force.
09:57
When you define your task, it has a contract
09:59
where you really have to define what
10:01
it means for the task to be done, given parameters.
10:04
Right?
10:04
It's not enough.
10:05
In the past, workflow tools, they
10:07
were considered done if it ran without exception.
10:09
And that's really not good enough anymore.
10:12
I'll get into kind of how we go around that.
10:15
And so we have a lot of work to do.
10:18
So the nice thing, too, about a lot of these frameworks,
10:21
right, this applies.
10:22
The workflow and kind of modeling your dependencies
10:25
in the DAG are completely decoupled
10:27
from what it's actually doing to run it, right?
10:30
We have Scala Spark jobs that run
10:32
to CP, Hadoop commands, that doesn't matter.
10:34
They're all codified in Luigi and Python.
10:38
So these Python, the Luigi tasks form
10:41
DAGs of dependencies naturally when
10:43
you wire your requires in Luigi.
10:45
We have parameters that we supply on the command line
10:48
when we launch these.
10:51
And then the parameters that, where necessary,
10:53
we propagate those down into the requires dependencies.
10:55
So I know some of this is standard Luigi stuff for people
10:58
who are used to working with it.
11:00
Most notable parameters are the label.
11:03
So when I say label, I mean music label, by the way.
11:06
And then dates, the active date that we're
11:08
running at for that label.
11:10
And this lets us reuse certain tasks across different labels.
11:14
And also, where tasks can be shared,
11:16
we just send the date if they're not label specific,
11:18
like the disk CP.
11:22
And then here's a visualization of that graph kind
11:26
of rotated 90 degrees so that it can fit on a PowerPoint slide.
11:31
Left to right, we're trying to actually run the push tasks.
11:35
That's what we actually launch.
11:36
We say, OK, it's a new day, push to S3, or Cron does.
11:40
And it's like, well, wait a minute.
11:41
Before I can do that, I have to have a single file.
11:44
OK, well, the single file doesn't exist.
11:46
I need to generate it.
11:46
And it just recurses all the way back
11:48
until it notices for the day that we're
11:50
trying to do that hasn't even transformed the data yet.
11:57
And so in a bit more detail, each of these pipeline tasks,
11:59
we have to codify three things, all of which are mandatory.
12:03
We have to codify any dependencies it has,
12:05
what upstream is necessary.
12:08
And we need a function to evaluate when we're done.
12:11
We need a way to actually say, I have arrived
12:14
at a state of completeness.
12:16
And then, most importantly, we need a function,
12:19
if it's not at that state of completeness,
12:20
to generate some artifact or external state
12:23
so that it now is at a state of completeness.
12:30
And here's what some of the code looks like.
12:31
This actually doesn't use the Luigi parameter API,
12:34
just because I thought it was easier to visualize it here
12:37
and just see it in raw Python.
12:38
But the push task takes in a date and label parameter,
12:42
because it's label and date specific.
12:45
And at first, it requires that there actually be a single file
12:48
that we've coalesced together with all the data into it.
12:51
And so that requires sending the date and label parameters
12:53
down the DAG into it.
12:56
And I kind of commented or wrote comments
12:58
for what the run and complete would do.
13:01
The run, if it's not actually done yet,
13:05
it means the run's going to execute.
13:07
And we just do the work of pushing that single file up
13:10
to an S3 bucket.
13:11
It's not terribly complicated.
13:13
And then complete, as you might guess,
13:15
is pretty straightforward.
13:16
The push task is complete if for a given label.
13:19
That artifact exists already up in S3 for the given date.
13:27
And then here's how the parameters actually
13:31
propagate down the DAG.
13:32
The push, coalesce, and generate stages
13:34
are all very label specific for the report we're building.
13:37
But we can, the discp and the transform steps
13:40
just need the date.
13:42
So those get reused.
13:46
And so the scheduler might be trying
13:48
to kick this off aggressively throughout the day.
13:50
Maybe the first attempt it gets, it recurses down.
13:53
It sees the transform isn't done, so it does that.
13:56
Then discp succeeds.
13:57
But then it gets to the point where it's generating.
13:59
It's running the Spark job.
14:00
And let's say there's some error or something,
14:02
some intermittent issue.
14:04
And it fails.
14:06
And we get some kind of alert for that.
14:08
But we don't freak out too much.
14:10
The scheduler, because this is all
14:12
scheduled to run automatically in CRAN
14:14
and we don't run these by hand, the second scheduler attempt
14:17
for the day can pick up right where it left off
14:20
because the discp tasks do the work of actually defining
14:23
what it means for them to be done.
14:25
So the DAG can kind of reuse work from earlier executions
14:29
because it just has functions to kind of figure out
14:31
where it was.
14:33
And this is kind of the core of the whole presentation about
14:36
the fault tolerance helped with the overall latency.
14:40
In teams I've been on in the past,
14:42
a failure like that, you'd have to redo everything.
14:45
Maybe the discp would have to be redone.
14:49
But yeah, we can use the scheduler to kind of
14:51
aggressively work off checkpoints.
14:55
And so this kind of bottom-up execution pattern
14:57
that everyone's familiar with, with Airflow and Luigi,
15:00
it leads to checkpointing your work.
15:04
It forces you to think about what decoupling points exist
15:07
and where you can reuse and where you can break off on
15:09
and say, at least I've done that for the day.
15:12
And you can, the most important thing
15:14
is you can schedule a daily task that you only need to run.
15:17
Once a day could be scheduled to run even as aggressively
15:19
as every 10 minutes.
15:22
And that's how our CRAN tabs start.
15:24
We still use Kronos, still a little bit older,
15:26
but that's how it's set up.
15:28
As soon as the new day rolls around,
15:29
CRAN is just automatically trying
15:31
to push the label reports for that day.
15:34
But it gets to a point where it can't because the data is not
15:38
there yet.
15:39
But as soon as it is available, it'll just launch.
15:41
And I'll go through that more.
15:42
We rely very heavily on the Luigi scheduler
15:44
to figure out what state the pipeline is in.
15:47
Is it ready to run?
15:49
Is it not ready yet because the data is not there yet?
15:52
Is it done already for the day?
15:53
We don't want to redo it if it is.
15:55
So just a simple thing is defining
15:58
that complete contract.
15:59
all this for free, which is the nice thing.
16:03
And again, it goes without saying,
16:04
if we're going to have a recursive bottom-up pipeline
16:07
like this, we need a base case for defining this execution.
16:10
And in our case, like put in English, our base case is,
16:13
we need a way to detect if the raw data that we're transforming
16:16
is available on the base pen or analytic clusters.
16:22
And our team, my team in particular,
16:23
has absolutely no control to force that, right?
16:26
It's not as simple as just writing a Luigi task,
16:28
and it'll execute that.
16:30
Those sources and all of that are completely
16:31
outside our means and our purview,
16:33
and we don't frankly want to own that.
16:35
So we can just codify what it means to have,
16:37
we can codify a task to kind of check for us
16:39
and then not blow up if it's not ready,
16:41
but kind of gracefully exit.
16:46
Leaf nodes in the DAG, in other words.
16:50
And this is our leaf node, per se.
16:51
We define it a little bit differently
16:53
than the traditional external tasks,
16:55
for those familiar with Luigi.
16:57
We use a bit of a different pattern.
16:59
It takes in a date and a minRows parameter,
17:01
which is used for validation.
17:03
It doesn't have any requirements,
17:04
doesn't require anything.
17:05
We're just going to run it right away if it's not done already.
17:08
And the job of it running is to transform
17:10
the events for the given day in the form that we want.
17:13
The raw events, all the spins data.
17:16
And we look at completeness as saying,
17:18
well, we're going to do a count over that new partition.
17:22
And if it's at least greater than some threshold, which
17:26
we've done the work offline to determine what makes sense,
17:31
then we'll know that it's valid.
17:33
We'll know that it succeeded.
17:34
So if the partition failed to generate,
17:36
the count will be zero.
17:38
It's not actually just blow up the query.
17:41
And even if there's some kind of malformed issue
17:43
with the partition, this is a very important defensive guard.
17:46
This has saved us several times, actually,
17:48
since I've been on the team, where
17:49
there was malformed data that came in for whatever reason
17:52
that I couldn't even dream up.
17:54
And we just had these in place.
17:56
And we've actually been in the position
17:58
of having to alert those teams because of these checks.
18:01
We can see, like, OK, something's wrong.
18:03
And we're the ones to actually tell the team.
18:06
I'm not trying to throw my other teams of Pandora
18:08
under the bus, but they have pretty good alerting.
18:10
But that's actually happened, where we were just
18:12
watching and noticed before them.
18:16
So the task is done if it passes validation,
18:18
if the min rows check succeeds.
18:21
We're pretty confident at that point.
18:23
And then again, these data delays happen frequently.
18:25
It's not always their fault. They
18:26
have complex syncs that have to deal with this.
18:28
And they manage the Kafka writes now, too.
18:32
But we can use the scheduler to not really care, right?
18:35
We're just going to aggressively schedule it
18:37
to run every x minutes.
18:38
I think it's defined as, like, 30 minutes in production for us.
18:42
And we supply the current date.
18:43
We say, every 30 minutes, just try to run this.
18:46
It's in human form.
18:48
It's like the, I think about it as, like, from The Simpsons.
18:50
Are we there yet?
18:51
Are we there yet?
18:51
Homer Simpson, while he's driving.
18:53
I didn't put a meme up for that.
18:55
But I guess I could have.
18:56
It's like, computers don't mind being asked that, right?
18:59
It's annoying for us.
19:01
But we can just constantly ask it, hey, are you ready to run?
19:04
Well, why aren't you done yet?
19:05
And just kind of annoy it.
19:07
And then, but the nice thing is, once the data becomes
19:09
available for transformation and that complete check passes,
19:13
we can launch the pipeline, right?
19:14
And it's completely automated.
19:18
So let's just go through.
19:20
I'll go through each stage.
19:21
Not code-wise, but I'll go through, kind of in English,
19:24
what each done state is for each of those five steps, right?
19:27
We talked about the transform one already.
19:30
There's a lot of Pandora Hive tables.
19:31
A lot.
19:32
But we want to narrow that down to just the spins
19:34
that we're interested in.
19:36
We create our own managed table on my team for that.
19:39
Partitioned by day, still.
19:40
But we do certain things to their raw data.
19:42
Like, we deduplicate them.
19:44
They do some of that, too.
19:45
But we do our own deduplication, just in case.
19:48
We filter out certain properties of songs.
19:50
Like, if it doesn't spin for longer than 30 seconds,
19:53
the labels seem uninterested in that.
19:54
Maybe I could get yelled at after this.
19:56
But I was told that, at least.
20:00
And then we have some other joins we do to kind of keep it
20:03
sane.
20:04
We keep the same partitioning structure.
20:06
We keep it by day.
20:06
And then we do other things for the metadata, too.
20:09
So that way I describe that metadata,
20:10
like listeners, stations, it doesn't start that way.
20:13
We have to kind of massage it down to that very natural form,
20:17
kind of denormalize it, so to speak.
20:19
So it's just easier to build this data set for.
20:22
And then, most noticeably, we just
20:24
don't want to be just CPing hundreds of tables
20:26
to do it on our side.
20:27
We'd rather take advantage of their cluster and use.
20:29
We use Hive on Tez.
20:30
It's very fast to build out some of those partitions.
20:34
And then we can just sync them back.
20:35
So if it passes the Minrose checks, we're done.
20:40
We consider that valid for the partitions.
20:42
We're actually considering switching this out
20:44
to just use snakebyte byte level checks for sanity,
20:48
because it's a little bit faster than having
20:50
to actually issue a query.
20:51
Even with Tez, it's still like we don't really.
20:53
We could be doing this in milliseconds
20:55
with snakebyte to actually check bytes for the part files.
20:58
We don't really need to issue a query.
21:00
So we would consider that if anyone was thinking that.
21:03
And we might do that just to make it better.
21:07
And then so once that's done, the disk CP stage
21:09
is ready to done.
21:10
We run an atomic disk CP, nothing really fancy.
21:12
We add it to the warehouse.
21:14
And very obviously, that stage is done.
21:17
If the partitions exist in Hive.
21:19
So we can very easily build out that state.
21:24
And then the generation stage kicks in.
21:27
We use Spark SQL for this.
21:28
So the legacy report was kind of a Hive MR job.
21:32
And it was a natural transition just
21:35
to take a lot of that with a few minor changes
21:37
and just run that right as Spark SQL over Hive
21:40
in a Hive session.
21:42
They're all parquet events.
21:44
So we get pretty fast filtering just from that.
21:47
When we're running the generation,
21:48
now we're getting label specific.
21:50
So we give it the task is going to run
21:52
with just a label on that control.
21:54
That's fed to the Scala job in terms
21:56
of how we're going to filter it down.
21:59
And we can actually do multiple label reports in parallel
22:02
this way too, which is what we do.
22:05
And then as part of that Spark SQL thing,
22:07
all the metadata gets joined in, artists, composers, stations.
22:11
We want to create a sync for every single spin
22:13
we want to make it as rich as possible
22:15
when we reported the labels for fields
22:18
that they've said they were interested in.
22:20
And then the most important part is in the driver directly,
22:23
we actually validate.
22:24
We do an additional validation on the number of rows
22:26
that come out of that stage.
22:28
And we have lower and upper bounds for this.
22:32
So it has to actually fall within an expected.
22:35
If there were maybe issues with the metadata
22:36
or something like that, even though we
22:38
had those defensive checks, maybe still something
22:40
went wrong.
22:40
You can't be too defensive.
22:42
So we check lower and upper bound limits.
22:44
Why do we do upper bound?
22:46
I guess we actually did have one case
22:48
where somehow new lines got introduced
22:50
into some of the fields.
22:52
And ever since then, we just don't take any chances.
22:54
Why not?
22:57
So if you want to throw off the music industry's reporting,
22:59
you could name your track or your band
23:02
Carriage Return New Line or something.
23:05
Not our reporting now.
23:07
So and then if it gets past this stage,
23:10
we output multiple part files to Hadoop, to HDFS.
23:16
We still do a lot of co-loc on-premise that we run HDFS.
23:20
We don't actually do a lot of our processing
23:22
like in S3 like that.
23:24
We have bare metal HDFS, if anyone was wondering.
23:27
I can answer questions about that later, for what I know.
23:32
So this is part of the Scholar job.
23:34
It's really kind of uninteresting,
23:36
but it just illustrates one point.
23:38
We run the Hive query.
23:40
We run it as just raw Spark SQL text.
23:43
And that text gets parameterized with string replace
23:45
before it gets sent to the session command.
23:48
And we still call, out of just habit,
23:50
like we call persist memory and disk,
23:52
even though cache, as of Spark 2,
23:54
is just the synonym for that.
23:56
If you look under the hood, we're
23:57
so used to cache from Spark 1x2.
23:59
is just calling storage level memory
24:02
and having to deal with that, that we're still
24:05
in this habit from our RDD days.
24:08
But once we have that data frame, the results,
24:11
we can call a .count on it, feed it to the check rowbounds
24:15
function, and then that will actually elevate an exception
24:17
if it fails the validation check.
24:19
And then the whole pipeline hard fails, and we get alerted.
24:23
But if it doesn't, if it passes, it just sort of silently
24:26
returns, and we can go on to just saving those results
24:28
to a label on date-specific HDFS location.
24:32
But the nice thing, too, it's a really nice use of Spark.
24:35
It kind of illustrates the caching benefit,
24:37
because we can do a validation step right in the Spark driver.
24:41
And because it's cached, it doesn't have to re-materialize
24:44
and re-evaluate everything to do the save.
24:46
I know this is Spark 101 stuff, but sometimes the benefits
24:49
of the caching aren't always obvious.
24:51
This is tremendously beneficial.
24:53
And if it fails validation, there's
24:54
nothing we have to clean up, right?
24:55
It's just an ephemeral data frame.
24:57
It just dies.
25:01
So then once all these part files are written out,
25:04
we need to coalesce them into a single file.
25:07
Some customers want compressed BZIP.
25:10
Others want GZIP.
25:11
I'm not sure which one is the orchard to check,
25:14
but we have to handle both of those cases.
25:18
And for BZIP, this is a little easier.
25:20
The output is already BZIPed part files,
25:24
and we can just binary cat those together into a single file.
25:27
Very effortlessly with Hadoop copy merge.
25:30
And then we get just a single BZIP file,
25:33
which there's even a nice flag to just clean up
25:35
the original files after it's done.
25:38
So if it succeeds atomically, it deletes those files
25:40
as part of that action, I think, atomically.
25:44
And then GZIP is a little bit more annoying,
25:46
but not too bad.
25:47
We just have to write out uncompressed files.
25:49
And they take up more space, but we
25:51
can use copy merge on uncompressed files
25:53
and do the same thing.
25:54
Wipe them out after there's a single compressed file.
25:57
And then just run that through GZIP,
25:58
and we get a single compressed file.
26:01
I put this here because I didn't always actually know.
26:04
I've been working with Hadoop for a while,
26:05
and I never even used this copy merge thing.
26:07
So if there's other people here that haven't used that
26:09
or didn't know about it, it's really handy.
26:14
And then the done state is this, predictably, right?
26:16
You might guess.
26:17
We expect a single compressed file in HDFS
26:21
for the label and the date.
26:23
And if it's there, we already consider it valid.
26:25
We validated all the data at the part file level in the job.
26:28
And we can operate off just that single file,
26:31
the assumptions around that.
26:35
So the S3 push stage kicks in.
26:38
And we just take the single file, put it up to a name.
26:40
The labels all get their own individual buckets, right,
26:43
for obvious reasons.
26:44
They don't want each other's.
26:45
They don't want to share their data with each other.
26:49
It makes sense.
26:49
So they get their own bucket, and then we
26:51
write out a new path in that bucket for the given date.
26:55
And then we can tell very quickly
26:56
if we're done with the S3 target classes.
26:58
And Luigi just, is it there or not?
27:04
So this is a short presentation.
27:05
But just the takeaways of this.
27:08
I guess if there's stuff to learn out of this,
27:12
we built all this so that it could be run very aggressively
27:14
from CRAN.
27:15
And it's a little confusing to see daily tasks that
27:18
operate with daily granularity being run every 10 minutes.
27:21
And it does produce more logs and Kibana and other things
27:24
that we have to deal with.
27:26
Finding the log that actually ran it can be annoying.
27:29
But it's a small price to pay for the self-healing aspect
27:32
and the automation that comes out of just,
27:35
we don't even have to think about how these are run.
27:37
If they fail, it's because there's a real issue
27:39
with the Pandora cluster or something,
27:41
or something we did with a patch that we put in trying
27:44
to make it better, and we screwed up.
27:46
But we don't have to guess when it's ready to run.
27:48
We just let the code figure that out for us.
27:50
And we just have CRAN try to just try over and over again,
27:52
annoyingly.
27:54
And if there's an issue, it just picks up right
27:56
at the most natural checkpoint where it succeeded.
28:00
Whoops.
28:02
So all of the tasks are compiled.
28:03
I sound like a broken record, I guess.
28:05
But the task composition is designed
28:07
around these checkpoints.
28:08
We had to think logically about where it made sense
28:10
to break these off.
28:12
Whereas the first iteration of the pipeline
28:14
didn't really do that.
28:15
It would build the single file wall.
28:17
It was generating the report.
28:19
And that had problems.
28:20
And it goes without saying, I've talked a lot
28:25
about the fault tolerance part of it,
28:26
because it's less obvious.
28:28
But using the latest frameworks, switching this over
28:31
to some of the latest, using Hive on Tez where we can,
28:34
and then Spark where we control it,
28:36
really helped for the parts that we can control,
28:38
using the latest open source frameworks,
28:40
just getting fast speed and execution.
28:42
So if there's any kind of issue with it,
28:44
it doesn't take forever to run or backfill.
28:46
We have to backfill more than we would like.
28:49
But that's it.
28:51
Are there any questions about?
28:52
I didn't cover how it improved.
28:54
I know I omitted that on purpose, or our environments,
28:57
or.
29:04
Yeah, just wondering, instead of using Crown,
29:07
do you have some sort of an event trigger
29:11
to trigger the pipeline?
29:15
We don't.
29:16
We could explore that.
29:17
That would make it easier.
29:20
I guess that would take the form of a task that
29:22
was polling the raw data.
29:24
But there still has to be an inception point to launch.
29:30
Yeah, I'd have to explore that more.
29:32
I mean, we certainly thought about changing this from,
29:34
it's right now, it's kind of a poll model.
29:36
And it could be more of a push.
29:38
But haven't explored that that much.
29:42
It'll probably, to be honest with you,
29:43
but probably before that, we'll switch it out to,
29:45
we'll be accumulating these directly off the Kafka streams.
29:49
And so the data sets will be generating throughout the day
29:53
as spins come in.
29:54
And then as soon as the day ends,
29:56
we'll be able to pipe it up to S3.
29:58
That's probably what will happen first.
30:03
So we use Luigi, and basically, to kick off PySpark jobs.
30:09
I'm just curious that you're using Scala Spark,
30:13
essentially.
30:14
How does that sort of work in terms
30:16
of calling Scala Spark from Luigi,
30:23
because Luigi's in Python?
30:25
Yeah, so in the Luigi contrib package,
30:28
there's a really nice Spark submit task.
30:30
It's actually Luigi contrib, and I think it's directly,
30:35
like Luigi contrib Spark is the import Spark submit task.
30:38
So you can use that.
30:40
It's really easy to submit a Scala Spark task.
30:42
It's a little trickier.
30:43
You can have PySpark tasks as well.
30:46
It's just a little bit trickier to use that.
30:49
I'll just say that the contributed,
30:50
they have a PySpark submit task that they've also contributed,
30:53
and that's junk.
30:54
We don't use that.
30:55
But it's the base one that's been contributed
30:58
makes it really easy to do this.
31:00
So you just point it to a jar, and you give it
31:02
the name of the class in that jar.
31:04
And the Luigi task's only job is to make a Spark 2 submit
31:07
command.
31:08
And so you can actually see, we print out that command.
31:10
We see what it looks like, and then it just runs.
31:13
So using Luigi is just an easier way
31:15
to not have to make a Spark 2 submit command by hand.
31:19
That's all it gains.
31:22
Other questions for Brian?
31:30
You mentioned you guys still use HDFS, like not relying on it.
31:35
Is there any reason for not crossing the line
31:39
or staying where it is?
31:44
Um, there is.
31:45
I mean, there have been political reasons in the past.
31:47
We do consider Amazon a competitor.
31:52
People use Amazon Music a lot.
31:53
I don't actually think those are the real reasons.
31:56
And I can say that because we're, the next two years,
31:59
I think there's a pretty.
31:59
aggressive plan to work out, to go into GCP, the Google Cloud.
32:04
We're going to transition a lot to that.
32:05
And they are still a competitor.
32:06
There's Google Music.
32:07
So I used to think that way for a while, and it's not that.
32:13
There's really no, I think it's just the bare metal.
32:15
We've always been kind of addicted to the bare metal
32:18
performance.
32:19
And we have really great sysops, like sysadmin people
32:23
that are just really good with Hadoop.
32:25
And if it got to be a problem where we would, we saw that,
32:29
we didn't have to worry about these issues.
32:31
But we don't really have a lot of those issues.
32:33
A lot of our hiccups come from the application level,
32:36
like the data engineers, like myself.
32:38
We screw up something in Hive.
32:41
We made a patch to the job that's
32:43
producing all the upstream raw data, and that failed.
32:46
So I think it has to do with the fact
32:49
that they just haven't been big woes.
32:52
They haven't bitten us yet in ways
32:53
that have become huge time sinks.
33:00
All right, any other questions for Brian?
33:04
Oh, we got one back here.
33:10
I was just wondering what you like about Hive,
33:12
that you still keep it around in addition to,
33:17
like I know there are some benefits,
33:19
but what makes it worth it to you
33:21
to support all that extra infrastructure over just using
33:26
Spark for everything?
33:30
Yeah, it's a good question.
33:32
Some of it's outside, not to defer too much.
33:35
There's a lot of teams at Pandora.
33:37
We're heavy Hortonworks users, and Hortonworks
33:39
has gotten a lot behind Hive in recent years with LLAP.
33:42
And so if you live in the whole Hortonworks environment,
33:46
you tend to kind of go down the pathways of those tools.
33:51
The Pandora, there's all sorts of different teams.
33:55
They're not that homogenous in terms of like,
33:57
this team uses Hive only.
33:58
There's all sorts of different activity on these clusters.
34:02
Spark, Hive on Tez, we have an Apollo cluster.
34:05
There's other warehouses.
34:08
There's very few Hive MapReduce jobs left.
34:12
What I'll say, what we've gotten addicted to
34:15
is kind of Hive as a catalog, as the meta store.
34:18
The meta store in Hive is probably
34:20
going to live on way past even many of these frameworks,
34:24
these execution engines that are down in the dirty with Hive.
34:29
We'll probably live on as a means of cataloging it.
34:32
And so that's kind of what Pandora has gotten addicted to.
34:35
Hive itself, in terms of the way we used to think about it,
34:38
like when I contributed to Hive, it's totally different now.
34:41
The people that use it in production
34:44
get sub-second latency with Hive now,
34:46
because they're running Presto over the meta store,
34:48
or Impala, or Tez, or something, or Spark SQL like we do.
34:53
It's just not too much of a concern,
34:55
and we get some nice properties out of it.