Modern Data Lakes Overview
developer.sh
developer.sh
Nothing about it was superior or even on par with simply fixing our current shortcomings OLAP database setup.
The data lake is not faster to write to; it’s definitely not faster to read from. Querying using Athena/etc was slow, painful to use, broke exceedingly often and would have resulted in us doing so much work stapling in schemas/etc that we would have been net better off to just do things properly from the start and use a database. The data lake also does not have better access semantics and our implementation has resulted in some of my teammates practically reinventing consistency from first principles. By hand. Except worse.
Save yourself from this pain: find the right database and figure out how to use it, don’t reinvent one from first principles.
I've found data lakes complement DW's (in databases) well. Keep the raw data in the lake and query as needed for discovery, and load it into structured tables as the business needs arise.
Data lakes alone are doomed to be failures.
(edit: the rationale behind this tends to be that you can avoid the heavy lifting of ETL/transformation logic by just using a data lake - obviously not the case, as most of us know)
There is after all a reason that the role Data Engineer became popular just as Data Lakes become popular.
S3 is basically free and has unlimited scalability. Oracle, DB2, HANA, SQL Server etc are ridiculously expensive and struggle under high concurrent load even with QoS in place.
None of the databases you listed there are OLAP databases.
Clickhouse, TiDB, Redshift, Snowflake, etc are significantly more suitable and should be the target of comparison here.
If you're able to solve the problems that you were previously using oracle or SQL Server for with S3, more power to you, but the truth is that to replicate the functionality of that old Oracle server you'll start with S3, but you'll also want some querying (Aurora? RDS? Hbase?), probably some analytics and ingestion (Redshift? Kinesis? Elastic? Hive? Oozie? Airflow?), along with some security now that you've got multiple tools interacting (Ranger? Knox?), probably some load balancing (Zookeeper?), maybe some lineage and data cataloging (Atlas?), etc.
In my experience what starts with "Just throw some data in S3, forget that old crusty expensive server!" ends with 22 technologies trying to cohesively exist because each one provides a small but necessary slice of your platform. Your organization will never be able to find one person who is an expert in all of these (on the contrary, you can find an Oracle, or DB2, or SQL Server expert for half the money) so you end up with seven folks who are each an expert in three of the 22 pieces you've cobbled together, but they all have slightly different ideas on how things should work together, so you end up with a barely functioning platform after a year's worth of work because you didn't want to just start with a $400k license from Oracle.
If you have S3 you can use Athena, Redshift Spectrum or Spark as query layer. It's not 22 technologies.
You don't need ElasticSearch, Ranger, Knox, Zookeeper etc as they have nothing to do with querying.
An OLAP database is, in the default case, an always-online instance or cluster, costing fixed monthly OpEx.
Whereas, if your goal in having that database is to do one query once a month based on a huge amount of data, then it will certainly be cheaper to have an analytical pipeline that is "offline" except when that query is running, with only the OLTP stage (something ingesting into S3; maybe even customers writing directly to your S3 bucket at their own Requester-Pays expense) online.
You'd be looking at $M in licenses for anything half-serious based in Oracle tech. Becoming good at replacing Oracle stuff probably has been one of the best paying jobs for a while.
My problem is the scalability and elasticity of it's licensing model. It doesn't meet the needs of today's analytics without spending enormous amounts of money up front.
That's why AWS has entire product suites from Athena, Redshift Spectrum, Data Lake Formation, Glue, etc to help companies actually do something with the files stored in S3. And it's often a mess compared to just fixing their processes and ingesting it properly into a SQL data warehouse first.
But data lakes have arisen from the enterprise where the centralised data warehouse was the standard for the last few decades. They know how to use a database. They know how to model and schema the data. And they know about all of the problems it has. They didn't buy into the data lake concept because it's trendy.
Fact is that for large enterprises and for those with problematic data sets e.g. telemetry databases simply don't scale. You will always have priority workloads e.g. reporting during which time users and non-priority ETL jobs come second. And often Data Science use cases are banned altogether.
The reason data lakes make sense is because it is effectively unlimited scalability. You can have as many crazy ETL jobs, inexperienced users, Data Scientists all reading/writing at the same time with no impact.
Generally you want a hybrid model. Databases for SQL users and data lake for everything else.
I do a mix of data science and software engineering, dealing with the datalake is a nightmare and I avoid it at almost all costs.
You know what the first thing everyone I worked with wanted to do after pointlessly pouring everything into the black hole that was the datalake? Re-implement some kind of SQL (and database semantics) back on top of it again; except now it's worse.
The data lake model seems to be more about not wanting to commit to a warehouse (for example: future proofing, looking at non-relational data, etc.).
Eh, almost all Data Lakes cannot handle small files well. All it takes is for someone to write 100 million of tiny files into the Data Lake to make life miserable for everyone else.
Every time I've seen someone do this it was a mistake and quickly resolved. Either you have way too many partitions in a Spark job or you are treating S3 like it's a queue. And if you really do need lots of delta records then just simply have a compaction job.
Nevertheless, my point is a Data Lake's does not offer free unlimited scalability. It takes a lots of effort and good engineering practice to make a Data Lake run smoothly at scale.
Once data has been through some transformations at the hands of a Data Scientist, it's now a secondary source—a report, usually—and exists in a form better suited to living in a Data Warehouse.
Data Lakes need a priesthood to guard their interface, like DBAs are for DBMSes. The difference being that DBAs need to guard against misarchitected read workloads, while the manager of a Data Lake doesn't need to worry about that. They only need to worry about people putting the wrong things (= secondary-source data) into the Data Lake in the first place.
In most Data Lakes I've seen, usually there are specific teams with write privilege to it, where "putting $foo in the Data Lake" is their whole job: researchers who write scrapers, data teams that buy datasets from partners and dump them in, etc. Nobody else in the company needs to write to the Data Lake, because nobody else has raw data; if your data already lives in a company RDBMS, you don't move it from there into the Data Lake to process it; you write your query to pull data from both.
An analogy: there is a city by a lake. The city has water treatment plants which turn lakewater into drinking water and pump it into the city water system. Let's say you want to do an analysis of the lake water, but you need the water more dilute (i.e. with fewer impurities) than the lake water itself is. What would you do: pump the city water supply into the lake until the whole lake is properly dilute? Or just take some lake water in a cup and pour some water from your tap into the cup, and repeat?
However, the scaling limitations of traditional RDBMS are insurmountable when trying to do things like data science, for instance.
It sounds very common sense to not to "limit the potential of intelligence by enforcing schema on Write" while in reality, the same problem just shifts (or gets hidden) in the next steps.
For example: there are 10 data sources with each 100TB of data. I aggregate these to my new shiny data lake with a fast adapter. Just suck it all without any worries about Schema. So, now I have 1PB of semi unstructured data.
How do I find the fields X and Y when these are all named differently in 10 sources? Can I even find it without having business domain experts for each data source? How do I keep things in sync when the structure of my data sources change (frequently)?
It seems like there is an underlying social/political problem that technology can't really fix.
Reminds me the quote: "There are only two hard things in Computer Science: cache invalidation and naming things."
and off by one errors!
The idea, though, is that, if your Data Warehouse wants the data in the form of e.g. a daily-aggregate accounting ledger, then your data sources might be of various time granularities and might be denormalized in different ways (one source with separate Invoices with Transactions foreign-keyed to an Invoice; another with just Transactions with root-level metadata like timestamp directly on them; etc.)
All of the transformations between the source formats and the destination format here are, in some sense, "transparent"—a sufficiently-advanced DBMS query planner could generate an OLAP expression to turn one into the other without understanding the problem domain. It's precisely because of this that, in many cases, it's cheaper to not worry about these kinds of transformations until you need to compute on the data. It's just a bunch of trivial stuff, that you can easily normalize in the computation step, but where fixing it on ingest would have been a whole expensive cluster operation to rewrite terabytes of data, and would require the OpEx of a whole additional set of always-online Hadoop cluster-nodes to fix marginal data as it comes in. Even though you're just going to be touching it all again anyway when you run it through the compute step.
If you are a relatively small operation, I’d recommend weighing additional complexity over the benefits. Sometimes a few we’ll written pages can suffice, other times you need to make the investment.
There is apparently now experimental support for using Hive partitions natively. Never used it, literally found out two minutes ago.
The number of records per object is usually "all of them" (restricted by partition keys). The main exception is live queries of compressed JSON or CSV data, because BigQuery can't parallelize them. But generally you trust the tool to handle workload distribution for you.
This works a little differently if you load the data into BigQuery instead of doing queries against data that lives in Cloud Storage. You can use partitioning and clustering columns to cut down on full-table scans.
The catch is if you need to filter by a property of the session, you are opening every session in range to check if it’s the one you want. That gets expensive quickly and is a bit slow.
For data lakes, parquet and Spark support fairly sane date partitioning. Partitioning by anything else is a question of whether you need it, such as a customer ID, etc. but remember this is a data lake, not a source table for your CEOs daily report. The purpose of the lake is to capture everything that you sanely can.
When you can’t store everything, usually due to cost, you then have to aggregate and only keep the most valuable data. For example in AdTech, real-time bidding usually involves a single ad request, hundreds of bid requests, a few bid responses and the winning bid. Value here is inversely related to size - bid requests without responses are useful for predicting whether you should even ask next time, but the winning bid + the runner up tell you a lot about the value of the ad request.
For structuring warehousing for reporting/ad hoc querying, to me the flatter the better - this uses the native capabilities of columnar stores and makes analysis a lot faster. Downside, good luck keeping everything consistent and up to date. Usually you end up just reprocessing everything each day/hour/whatever the need is, and at a certain point say no new updates to rows older than X.
The cool thing about modern data warehouses, is that they include interfaces to talk to the data lakes, so your analysts don’t have to jump to different tool chains, such as Redshift Spectrum (which is basically Athena) and the aforementioned BigQuery ability to use tables, streams and files from GCP.
It’s an incredibly productive time to be working with all this! Even 10 years ago, you’d need a lot of budget and a team to just keep the lights on, today it’s all compressed into these services and software.
Sounds like you are thinking more of a data warehouse, which is structured data on an engine that’s designed for querying large volumes of data. I’d recommend first starting with your objectives and then going for what solves with least amount of “stuff”.
I don’t work on data warehousing or pipelines now, but when I did a year ago, AWS and GCP both offered great tools with slight differences, where AWS was a bit pricier to start, but focused on more predictable pricing and GCP was much cheaper with pay as you go, but you could get yourself in trouble with cost by not following their best practices.
I’d recommend starting off with an OLAP database and going from there, reaching for a datalake once-and only once-you’ve reached the limits of the OLAP db.
I've gotten pretty far with jsonl on gcs and bigquery - even some bigquery streaming for more real-time stuff.
https://cloud.google.com/bigquery/external-data-sources
My biggest papercut with using this was having to make sure that all of the locations matched exactly.
My comment was intended for those just starting out. If you don’t really know what you are doing yet with data, it best to focus on your core company objectives and not burn valuable engineering time on infra you can buy for now. Unless that data stack is your core business.
- The cost of storage is the same as S3.
- Storage and compute can be scaled independently.
- You can store multiple levels of curation in the same system: a normalized schema that reflects the source, alongside a dimensional schema that has been thoroughly ETL’d.
- Compute can be scaled horizontally to basically any level of parallelism you desire.
Given these facts, it is unclear what rationale still exists for data lakes. The only remaining major advantage of a data lake is that you aren’t subject to as much vendor lock-in.
You can save plenty of money if you have the scale to move out of S3. That’s important because you can usually trade CPU for storage by storing data in multiple formats, optimized for different access patterns.
But mostly, the Hadoop ecosystem is very open. The tools are still maturing and it’s easier to debug open source tools than dealing with the generally poor support in most managed solutions.
Why does it take scale to move out of S3? And I thought S3 was cheap, so how would moving out save money?
But to be fair, you'll go on-premise due to the computing or bandwidth costs first. And you'll likely move data to the same datacenter to avoid expensive transfer costs.
I've also had to work in places where you simply could not put your data in the cloud due to regulatory reasons.
https://medium.com/@vtereshko/data-warehouse-storage-or-a-da...
(Pm on BigQuery)
There is presently strong interest in associating this data with other DBs, of which I am aware of about 80, with a total of probably 500-1000 tables, along with some very old "nosql" b-tree datastores in MUMPS. There are new $10M+ projects coming online around the enterprise roughly every day.
Where would you start?
It's very easy to setup so you should be able to test it quickly to see if it fits your needs.
Microsoft SQL Server with Clustered ColumnStore tables would make practically all queries fast on that, especially if most queries are only for subsets of the data. PostegreSQL could probably handle that too, no sweat.
Also see "Your data fits in RAM": https://news.ycombinator.com/item?id=9581862 which would mean that you could do in-memory analytics of your relational data with SAP HANA or SQL Server if you really needed that kind of performance: https://docs.microsoft.com/en-us/archive/blogs/sqlserverstor...
You can spin up either SQL or HANA in the cloud or on Linux, so you don't even need Windows. Both can be connected to just about any other database you can name, often directly for cross-database queries. SQL 2019 is particularly good at virtualizing external data: https://docs.microsoft.com/en-us/sql/relational-databases/po...
10 gigapixel images are a completely separate problem. If you need individual images to be fast to view, you want some sort of hierarchical tiling like Google Maps does. If you're processing them with machine vision or something, then you want whatever makes the ML guys happy.
PS: I hope you're not working on DARPA's spy drone, because then please disregard everything I said and delete your data for the good of humanity: https://www.extremetech.com/extreme/146909-darpa-shows-off-1...
My question is more the specific mix of problems: a DB, a ton of image data, and other adjacent DBs that people want us to play with. How would you set that up?
I work on cancer, so, definitely not spy drones.
For joining data from multiple database, if the data is large, I would use something like Presto(https://prestosql.io/) to join and process the data. But that's partly because we have already had Presto clusters running.
Similarly, it's actually hard to beat MS SQL Server for OLTP workloads, especially at moderate (~1TB) scale or for ad-hoc queries that require parallelisation but not distribution to a cluster. In other words, it's great for "Medium Data".
It does actually scale to large clusters with the new SQL Parallel Data Warehouse: https://docs.microsoft.com/en-us/sql/analytics-platform-syst...
That's also available as an Azure service if you want to have a play: https://docs.microsoft.com/en-us/azure/sql-data-warehouse/sq...
But realistically, distributed clusters are almost certainly not what you need. They're complex and slower for simple queries that could be answered by one box with a good indexing scheme. Just to reiterate: for large tables with hundreds of millions of rows, you want a modern, column-oriented database. I can't stress this enough: if you haven't yet played with SQL's ColumnStore, go spin up an instance in Azure or AWS and give it a go on one of your larger tables. It's crazy good. I've seen compression ratios of 50:1 and query performance improvements of 300:1 with basically zero hand-tuning of indexes or any such thing.
There's a reason people pay $10k+ per core for Enterprise SQL Server licensing. But hey, if you're penny-pinching on a $10M project, then as I said, MySQL and PostgreSQL will work. They're better at replication, clustering, and MySQL (only) is better at low-latency for trivial queries. But they tend to be poor at connecting to commercial or otherwise quirky data sources. So then you'd probably have to layer something like Apache Drill on top: https://en.wikipedia.org/wiki/Apache_Drill
I will never trust a company that stores everything in one data lake, that's major data breach just waiting to happen.
As an example, Athena would do a terrible job at finding a specific user by its ID, while Spanner would behave just as poorly at calculating the cumulative sales of all products for a given category, grouped by store location (assuming many millions of rows representing sales).
Hope this analogy makes sense.
The only "newsql" database that truly does OLAP+OLTP (now called HTAP) well is MemSQL with it's in-memory rowstores and disk-based columnstores.
Welcome to try it in March with TiDB 3.1.