Making Netflix's Data Infrastructure Cost-Effective
netflixtechblog.com
netflixtechblog.com
Stick to simple systems like elastic beanstalk or app engine + RDS and something like snowflake and you get powerful, manageable, auditable projects. I'm also finding great benefit in launching each project in its own AWS account. This means you know exactly what project costs how much!
Need web hosting? Certainly you can't deal with a 2Tbps DDoS attack, so let's route that traffic through Cloudflare.
Sending email? Can't send it out with those /23s you have, they're on some blacklist from a shady DNSBL provider. Need to pay for GSuite or Office365.
Can't protect your CEO against a targeted attack? GSuite Advanced Protection or Microsoft ATP has you covered.
These are really just variations on the theme of economies of scale, aren't they?
If your production system is 10 instances on some random cloud, a 10% efficiency savings saves me 1 instance, so maybe $2k a year. Taking into account opportunity costs vs doing things that raise revenue, said startup would consider the effort a waste of time unless it took a few hours.
With the same architecture, but instead 20K instances, then suddenly that 10% is saving 2000 instances, and 4 million a year. Unless there's a major engineering shortage, chances are that spending a over a month on 4 million in yearly savings will be completely justified, and would even be a highlight in someone SRE's review.
Chances are that the optimization wasn't even any harder in the big tech company: It's just that small savings on big piles of money are suddenly worth it. It's not possible to find millions of dollars in loose change between the proverbial couch cushions if you didn't have the chance to spend hundreds of millions in the first place. Heck, in a growth company, even at that size, 4 million might not be enough savings, as there might be even bigger things.
Same with dev tooling: Saving a company 5% CI times is not great with 10 developers, but if you are making five thousand developers more productive, suddenly you can hire an entire team of very specialized developers, and it's not a luxury.
This is also why often copying what large companies is foolish when you are small: The tradeoffs are going to be completely different. What would be an unacceptable flaw in a large company is just fine in your startup.
The real trick IMO is the 200-800 range: Large enough that the simple solutions for small companies have probably broken down and are causing pain, but nowhere near the staffing to hire yourself a team to, say, add types to Ruby, get a team to build around the weaknesses of your database, or whichever other problem you have that any member of FAANG would just throw 10 million dollars in staffing costs without batting an eye. I've seen way too many companies that stopped being able to grow their company at those intermediate sizes, and get stuck in technical hell.
I think the most applicable in our field is "resource cost accounting", which looks superficially like what Netflix described. You keep "resource accounts" in the units of the actual stuff consumed by your business processes, then map backwards from the final product to costs.
Of course I am not an accountant. Just took a course and found it enlightening.
Snowflake, Hive, Cassandra, RDS, Druid, Presto
Is there a good reason for this that people without experience at FAANG orgs can't grasp? To the naive, it seems like mixing Postgres, MySQL, Mongo, and Oracle.
Also why is it on Medium under a Paywall, they are worth, quite literally, nearly $200 billion dollars =/
Snowflake is a Data Warehouse/Lake, but it's also it's own custom SQL DB.
Cassandra is a NoSQL DB
RDS is running (some standard relational DB)
Apache Druid is a columnar analytics DB centered around realtime uses. It has it's own query language and delegates to Apache Calcite for specific DB/datasource drivers. Can integrate with Kafka/Hadoop.
Presto is (to my understanding) like a meta-DB that can query multiple databases. Similar to an integrated Apache Calcite or Google's ZetaSQL.
There are a LOT of overlapping concerns here, which is why the confusion.
Essentially they have 2-3 different products in several categories targeting generally the same usecase.
Everything else solves different problems at Netflix scale. I'm spitballing here but Cassandra could be used for metadata serving (high throughtput, embarassingly parallel reads with high uptime), RDS for their billing system (transactions, ACID, etc), Druid for realtime OLAP and Presto as an interface to Hive/Snowflake.
A smaller company wouldn't need this level of complexity (if you aren't large, you could probably serve your metadata from MySQL, and just use Snowflake as your OLAP engine).
For others reading your comment though, I did want to list some things I have used with Postgres that relate to connection pooling and data partitioning:
* PGBouncer for connection pooling/sharing.
* Postgres Table Inheritance for table partitioning (https://www.postgresql.org/docs/12/ddl-inherit.html)
* PGPartman for automating the creation of partitions (https://github.com/pgpartman/pg_partman)
* Citus for low-barrier data sharding (it's a Postgres Extension like PostGIS) (https://www.citusdata.com/)
I think yes. It's about the scale of teams working on different parts of the products.
Its not like one team runs a DB where all other teams all put their data into. So multiple team make their own choice of DB to use for their own use cases, and with their own preferences and reasons. So as a whole, Netflix or other FAANGS can grow to use many DB, many languages, many frameworks, etc.
They also have billions to tens of billions of dollars worth of hardware costs per year so what can seem like a waste of time to optimize can end up saving millions of dollars worth of hardware a year.
Maybe their scripts determined them hosting on Medium was costing them too much money. /s
Others have already mentioned several reasons that would justify their choices. I completely understand why a company of Netflix's size has half a dozen different databases.
Recently, I ran into a somewhat similar situation that left me asking "why the hell are there five different databases running on this ONE virtual machine?".
It was one of VMware's virtual appliances running one of their products. I want to say it was either vRealize Network Insight or vRealize Operations Manager but I'm not 100% sure. It was likely one of those or something along those lines.
To stand up the smallest possible deployment of this virtual appliance, the "minimum recommendations" were something like 8 or 10 vCPUs, 32 GB of RAM, and 800 GB of storage -- and, if memory serves, this was for a single node!
Anyways, I got looking into it a bit and discovered that there were a total of five (5) separate database instances running on this one single virtual machine. I can't recall the specifics now but I'm pretty sure there were two PostgreSQL instances and one MongoDB instance. The other two have escaped my memory but I believe #4 was Cassandra. NFI now what #5 was.
On a positive note, afterwards I completely understood why the "minimum recommendations" were what they were. I would have (naively?) thought it'd be better to spread those database instances across a few VMs but I suppose (from VMware's perspective) when your customers call up screaming because their appliance they just paid six figures for is running like crap before they've even really put it into production, 1) you've got just one VM you've got to deal with and 2) the easy "fix" is to simply tell the customer to give it more resources ("another 10 vCPUs and an additional 64 GB of RAM and it should be fine!").
</rant>
I collected a list of these starting at my previous employer and got up to 91 examples. Stuff like "consistent configuration names" (ssl_verify vs verify_ssl), "consistent log format", "consistent Kubernetes CRD naming scheme" etc are all examples.
It is not an easy thing to solve. Especially since folks don't usually think of it in economics terms.
Disclosure: I work for VMware, but not on the products you mentioned.
Better distribution. Medium will recommend the article to their readers.
You need extreme diversity of these tools to benefit from inventiveness of your staff.
Company mandates to strictly limit to one way will just beat creativity out of people and burn them out.
This is especially stupid if the reason for enforcing just one way or just one tool is to make SRE lives easier. Aside from that being the tail wagging the dog and making no business sense, you have the power to use software abstractions to still make SRE lives easier even without limiting freedom to use any system, even in a small company.
Engineering isn't art. It's engineering, a field where creativity and innovation works within economic, technical and moral constraints.
of the analytics tools some of these are storage systems, some are compute engines, and some are both.
snowflake is both a storage system and a compute engine. it is very nice but it is also very expensive. there are a lot of computations or analyses that could be ROI positive but that aren't ROI positive on snowflake.
hive has two main components, hiveserver2 and the hive metastore. hiveserver2, which is the part that executes hive sql queries, is a compute engine. the hive metastore is a storage system (in combination with hdfs or s3 typically). hiveserver2 has mostly been superseded by spark sql and presto at this point although it's still used in a lot of places for legacy reasons. the hive metastore is a critical component of a data lake as spark, presto, and most other standalone compute engines rely on it to determine where the files are for a particular dataset, how to read them, and what the schema is.
presto is a compute engine that executes sql against data typically stored in s3 or hdfs in parquet or orc format and registered in the hive metastore. it is excellent for low latency querying but all query execution state must be able to be held in memory. for queries where the state is too large to fit in memory, spark is a better fit. presto can also read data from other datastores such as mysql but this is a really good way to break your website (see point above about isolation) and is best avoided.
druid is both a storage system and a compute engine and use cases overlap heavily with presto. druid does two things that presto does not. the first is that druid allows you to update your data in near realtime. this is extremely challenging to do in hive/s3/hdfs based systems. the second is that druid maintains what is essentially a search index on top of your data which in some cases can dramatically improve performance of filtering operations. unfortunately druid stores all its data on locally attached disks and not s3, so data storage costs in druid are around 10-15x higher than a presto/hive/s3 solution. regarding near realtime updates, apache hudi and databricks delta lake are bridging this gap for presto/hive/s3, so unless druid can separate storage from compute and be used as a standalone compute engine i don't see it surviving in the long term.
I don't really care for personalized recommendations. As long as the movies and TV Shows have some type of keyword or grouping, I am happy. Given the relatively fixed size of the catalog, and that most movies and TV shows are already classified (for example as Action or Romance or Comedy etc) this doesn't seem like something where a lot of data would do a lot to help.
For me, most streaming movie companies are pretty equivalent in their experience. The only data I really need them to keep track of about me is what I have watched so I know if I have already watched something before, and where I am in the movie/tv series, so I can pick back up when I log in again. Everyone pretty much does it. As for streaming quality, basically everyone knows how to transcode movies and put them on a CDN. The size and quality of the catalog is the biggest differentiator for me.
Absurd.
It's a vanity URL for Medium. If they've got their own domain why would they want to drive any traffic to medium? That doesn't make any sense to me.
https://datacenterfrontier.com/mapping-netflix-content-deliv... [2016] https://openconnect.netflix.com/