Uber’s Big Data Platform: 100+ Petabytes with Minute Latency
eng.uber.com
eng.uber.com
Also interesting: A lot of companies, you look at their big data ecosystem and it's littered with tons of tools. Uber seems like they've always done a good job keeping that pared down, which indicates to me that their team knows what they're doing for sure.
The proliferation of tools allows to test out different ideas in production conditions which is not bad per se. But the duplication of efforts or variability has to come with some advantage to continue doing it. Infrastructure engineers wanting to reinvent the wheel can create a lot of technical debt fast, and leveraged one at that.
This takes nothing away from the great work Uber has been doing in the area. They have definitely set a great direction for their platform.
Consider Netflix. They are big open source contributors to the data industry. Many of their tools interest me because I can see their value outside of Netflix. Uber on the other hand is solving esoteric problems with esoteric solutions.
Having a clearly defined schema that can be shared between teams (we had a specific repo for all protobuf definitions with enforced pull requests) significantly reduces the amount of headaches down the road.
1 billion bytes = 1GB,
1 billion KB = 1TB,
1 billion MB = 1PB,
10 billion * 10MB = 100PB.
Which again leaves me wondering why Uber needs 100PB.
And it's not like Big Tech has ever misled people on this subject.
Remember, Uber also stores application positioning telemetry as well, so it isn't just the number of vehicles.
I had a friend head off in 2011 to work with Hailo in London, and designed much of the data ingest & storage architecture there. That was built on Cassandra, and some of the numbers were staggering, even for a (initially) single-city, < 2,000 driver system.
You're not just tracking start and end point of successful rides -- you're tracking users when their app is open (both drivers and users move around before and after initiating a ride, and you need to share that data between parties), you're also tracking empty vehicles so you can identify over / under serviced areas ... doubtless lots of other real time stuff you want to know. But the big thing is you want to then have that data available for mining, improving your service, regulatory obligations, and so on.
Most of what you describe has little meaning or use at a few months old. Sure you might want to collect detailed awareness of individual drivers over say a month, generic traffic patterns over six months, and little or nothing beyond.
Once the data has outlived that direct, primary, purpose you dispose of it or it's a famous, excessive, data breach waiting to happen. There is absolutely nothing wrong with only keeping billing data beyond say six months max.
The really invasive pattern-matches happen best retroactively at that scale.
Age of the data does not typically impact many such analyses.
There's a huge spectrum between your 'at most 30d' claim in another reply, and 'all data in perpetuity'.
30d is way too short. Perpetuity is probably the idealistic intent, and depending on storage costs over time this may be achievable.
> I suppose it demonstrates the difference in US and European perspectives on data.
How so? I described a the case of one UK-based company, and TFA is about one US-based company, both doing similar things.
> Most of what you describe has little meaning or use at a few months old.
Clearly this is not the case, as evinced by the actions of the people designing, building, operating, and maintaining these types of systems.
That said, 100kB does seem high. 10kB or so would seem more appropriate, as a very rough ballpark.
I've experimented with Parquet data on S3 for a work POC, and the latency to fetch the data/create tables/run the Spark-SQL query (running on EMR cluster) was quite noticeable. I was advised that EMR-FS would make it run quicker, but never got around to playing with that. But I guess the creating of in-memory tables using raw data snapshots would still remain true? Or maybe I missed something.
Also, I take it if 24 hrs is the latency requirements for ingestion to availability of this data, obviously this isn't the data platform that is powering the real time booking/sharing of Uber rides. I'd be curious to see what is the data pipeline that powers that for Uber.
Disclosure: am CEO of Fivetran.
Organizations that have the resources tend to go open source with a lot of custom tools (such as Uber, or FAANG). Those companies tend to be the "makers" of such tech as well.
Organizations that don't have those resources or in-house experience, or desire to build, rely on the commercial licensing for open source offerings, cloud based or traditional commercial offerings.
At these scales of data and time horizons, owning the engineering resources capable of supporting the tools might be necessary, but surely expensive.
If one can go with a stable, well known commercial offering that reaches the desired scale but keeps the volumes of data in open standard formats, I think that is a good compromise. Many commercial vendors have gone that way, for example, Microsoft recently went all in on integrating Spark and HDFS closely with SQL Server 2019, and a lot of other database vendors have already done that as well (e.g. HP with Vertica on Hadoop, etc.).
I also think it's possible that at these scales and performance, people doing that big data work really are on the bleeding edge and therefore, innovation and new development is almost a requirement to make it all work together and meet the aggressive performance and efficiency benchmarks desired. Especially when having to do complex things like handling both the speed and batch layers in one unifying architecture (as in lambda architecture).
For example, imagine that System X can scale well up to 1 terabyte of data, but that after that it becomes increasingly complex to handle more data and it requires special optimizations to maintain an acceptable level of performance.
From System X's perspective as a vendor, 99% of their customers are perfectly fine with the performance of System X and just want more features (that the vendor can charge for upgrading to the next version of System X). On the other side, there's that annoying 1% of customers just keeps on asking for more and more performance optimizations that 99% of their customers don't care about.
From the vendors' perspective, it makes more sense to invest development efforts into features (that can translate into more money) than in performance optimizations that only appeal to a narrow segment of their users. For the company that requires these optimizations, since the commercial codebase is proprietary, they're locked in and have no easy way out.
That's pretty much why the large tech companies invest in owning their own part of the stack. It mitigates a scaling risk, keeps expertise in house and ensures that the solution that is developed in house matches 100% with the often specialized needs of the business.
That being said, we use Snowflake heavily at Strava and are very happy with it.
That being said, there is definitely a large amount of NIH promotion-ware being developed at large tech companies, as you mention.
The pricing is for full scan so it is not just for in-place update but charged with bytes scanned for selected columns for scanned partitions.
If you invest a chunk of that money saved by not having to do 3 or 4 way data replication into good network, a lot of that move-compute-to-data complexity becomes unnecessary. The second big efficiency promise of the basic map reduce is doing lots of sequential I/O on those pesky disks which can do 100x sequential throughput compared to their pithy random read performance. Alas, in today's mixed workloads, you're certainly not getting those >100MB/s from a disk that might make a benchmark scream. So the network savings (you're still going to do that shuffle before your reduce anyway, aren't you?) from compute-next-to-storage becomes less useful.
Best I can tell, the compute-next-to-storage trick does still matter as you get to large data sets (PB in a job), but then as you keep growing your infrastructure (and I'd wager Uber would be just about at this point), the benefits of disaggregation start weighing increasingly heavily. In particular, fleet management becomes considerably easier with less resource stranding.
Just list a few: namenode HA, unbalanced r/w (a lot of old data but less used, wasting cpus for just making the machines up), cross dc replication.
You may be able to get some of them from cloudera, but would that be cheap?
Also, I assume that they have dashboards that use pre-aggregated tables for faster results, they probably have ETL jobs for this use-case but is the pre-aggregated data stored on HDFS as well?
I guess they want to avoid liability as much as possible?
Would it really be possible to use a blog post in a legal proceeding to determine whether Uber has drivers or partners?
Also, a lot of countries have laws saying [taxi] drivers _have_ to be employees, irrespective of VISA stuff.
“Big” data was never about how to store the data.
There is no size requirement. It's more about what you collect, the frequency, the coverage, and how you use it.
And anything you could fit into, basically, one well cooled room is not "Big" anymore, sorry to tell you that
100TB SSD from Nimbus, https://nimbusdata.com/products/exadrive-platform/scalable-s...
My point is, you could fit all those Uber data (I know, I know, replication, sync etc) into racks in the SINGLE well cooled room in datacenter, managed by 2 guys per shift.
And this is not "Big" as far as I can see.
Probably will speed up their queris/analytics as well all things being local