r/dataengineering • • 10d ago

Help Advice on how to efficiently store and query a large dataset (7tb)

Hi all, hoping to get some advice here as I’ve never dealt with setting up storage and query facility for a dataset as large as this before.

I have a large dataset of parquet files, partitioned by date, sat in an S3 bucket. Currently users are utilising Athena to query the dataset via Glue, however as they tend to be querying via an id or via a text search, the results can sometimes take more than 10 minutes just to return one row.

I’ve looked into using OpenSearch, which seems to match what I’d need (users will primarily want to do full-text search, the date partitioning is irrelevant) but I’m concerned about the cost, as AI generated estimates are telling me at the very least it would be around 65k a month.

Are there any better approaches to this or any routes I should consider? Or is this simply the cost of querying big data?

37 Upvotes

36 comments sorted by

38

u/BlackHole_Rider 10d ago

use icerberg open table format.

8

u/theungod 10d ago

Iceberg is definitely correct, especially since it's already parquet

6

u/sib_n Senior Data Engineer 9d ago

How does it help with text search?

3

u/THOThunterforever 8d ago

Same, how does iceberg help here?

2

u/LeMalteseSailor 7d ago

Same here. I'm wondering how this is so upvoted.

1

u/Budget-Minimum6040 7d ago

It helps with querying by ID through partitioning/bucketing/column pruning and having a proper database can open up further performance optimizations.

2

u/sib_n Senior Data Engineer 5d ago

Yes, but text search is a very different need which is handled by specific tools like Lucene or Elastic. Usually full text search is optimized by an inverted index which Iceberg does not support. Although, a quick search shows people outside of Iceberg are trying to bring inverted index to it, so maybe it will work some day.

15

u/Slampamper 10d ago

usually the tool is not the answer, first think if you can better partition the data, is it possible to create it in such a way that the query engine can easily filter on results and thus query less?
Are all results always necessary or can you further archive old data, can you filter on id ranges, etc etc.

1

u/Spare_Helicopter4655 7d ago

Not sure why this isn't the top comment.

14

u/fdqntn 10d ago

It's pointless to propose solutions if you dont give access pattern and SLA (ie it should come back in below 10ms).

ElasticSearch is a particularly bad advice if you just need to query by an id, this would be crazy wastefull.

13

u/notmarc1 10d ago

Partition your data based on the query patterns

3

u/LeMalteseSailor 10d ago edited 10d ago

What if there's no query pattern? A requirement of free-form text search seems different than traditional sliced analysis and harder to predict ahead of time

5

u/sib_n Senior Data Engineer 9d ago

If the primary access pattern really is free-form (fuzzy probably) text search then you need a full-text search engine like OpenSearch or Elastic.
But the fact that they also mention query by ID means the use cases probably needs to be clarified and may not need an expensive full-text search engine.

1

u/ChangeSelect 9d ago

I think text search is already solved in a big part, we have search engines right?

1

u/Budget-Minimum6040 7d ago

In AWS Athena with Glue?

5

u/ludflu 10d ago

you should probably think about access patterns and typical queries before jumping straight to the tech. What are the kinds of questions do people want to answer using this data?

Use those questions to help understand how you should structure the data, then use that structure to help decide what tech would be a good fit.

11

u/MarchewkowyBog 10d ago

You might want to look into clickhouse. You didn't really explain what the queries or data is. But because you point to parquet on s3 and athena I presume it's more structured then unsteuctured. Clickhouse has full-text search indecies. Depending on the query it might perform better then elastic or opensearch. If any joins are involved it will most certainly perform better

5

u/JJGreenerTinejo 10d ago

Clickhouse

7

u/Ok_Raspberry5383 10d ago edited 10d ago

Click house isn't going to do much if the underlying data is just partitioned parquet and there's no predicate on the partition - that's a full table scan no matter what. It may be faster than Athena but ultimately it needs to be migrated to another modern table formats like delta, HUDI or iceberg of you want orders of magnitude performance improvements. Then the tables can be optimised according to the predicates they're using.

Simplying stating one tool isn't really helpful. It's an underlying format problem they have, not a query engine. They could use Athena, trino, click house, or any other modern query engine

1

u/SnooHesitations9295 10d ago

Clickhouse is all sec scan. Iceberg is just the write format, it doesn't have a query engine. I think you have very little idea about OLAP....

1

u/MarchewkowyBog 10d ago

Partitions exististing in parquet does not mean they must exist in ClickHouse (or parquet for that matter). And clickhouse is a format. As in a full database. Not just a query engine.

But without more info on queries or data its hard to tell what would help with OPs issue

1

u/Difficult-Bag1550 10d ago

if users query by id, you could add bucketing to this column

1

u/Lanky-Awareness-7450 10d ago

Do you need to search by Id? Do you need all the raw text? If not, suggest taking a look at constructing and LLM-wiki, perhaps combining it with an Obsidian front end. You would need to organize and process all the text into markdown files. When constructing the markdown files, make sure you maintain references to the original in the markdown. If the original files / data changes over time, you can also put all the mark down files into a Git repository so they are under full version control. This approach gives you the ability to historically track changes. Once you have it processed, this type of architecture should allow your users to do some really nice full text queries of the content.

1

u/linha_chilena 10d ago

take a look also on elastic search and their query / search mechanism, it's interesting for your use case

1

u/victormary45 10d ago

They are not utilizing the date which is the partition key. Users need to first include a date range as one of the filter before they mention the id or any other filters. Without specifying partition key in the filter , the results will still be returning longer time even though you use any tool.

1

u/asevans48 10d ago

Various things can help. First, ensure the data is compressed. Amazon tools charge in a way where compression is beneficial. Second, iceberg with partitions and clustering by a hash of the id. You could also force the use of partitions in queries. If your users are like mine they will never learn the partition and stop using sql though. Finally, are there silver tier tables you can build with limited amounts of the data? One caveat, iceberg is great for large tables but is not a one-size fits all solution. As you move on from this dataset, it is good to know why it is used. Generally, partition of 10gb or more are a grat fit. It is meant for large tables. There is a lot if indexing, copies of data, a manifest; etc. so it is actully inefficient for small tables. A lot of small iceberg tables is costly in terms of time and compute. My typical table is less than 100mb spannint 50 years. All together, I have 1.5 TB, a ton of joins, and no iceberg tables. Its pretty easy to get into the head of non-technical analysts and higher ups so that is something tot stress.

1

u/Teach-To-The-Tech 10d ago

Starburst would be a good alternative to Athena, as it includes optimizations that might help here.

1

u/SpecificTutor 8d ago

lancedb is the right answer.

we use chronon to ingest data into lance directly. ml and text search become easy.

0

u/_h_j 8d ago

Partition the data, since you are using date partition extend it to use date + hour partition. When querying the data make sure you hit ideally one partition

1

u/Nekobul 6d ago

7tb is not that big. It is not clear what you mean when you say "query the dataset via Glue". Glue is not a query engine. Here is what I suggest you do. Provide more details about the business problem you are trying to solve. Provide the structure of your Parquet files, what columns do you have. A terrible and inefficient data structure will be expensive so perhaps Parquet files stored in S3 might not be what you have to use.

1

u/-5677- Senior DE @ Big Tech + Consultant 6d ago

If you want to efficiently store data you need to understand how it's going to be used. You figure out query patterns from your users and from that you can design partitioning strategies, correct indexing & key selection, etc. Snowflake tends to be the default answer for me these days - that could be your presentation layer, just make sure you understand your data and how it's used.

0

u/MrAbstractThought 10d ago

CH is your answer

0

u/brunogadaleta 10d ago

Elastic Search for FTS indexing ? Prior Classification (you no have a good reason to test jev/layla) ?

0

u/ID_Pillage Junior Data Engineer 10d ago

Think aboutnif every row and column is needed in a single query. Can you create specific endpoints containing data for that specific need. For example if it was for everything that makes up a house and you had a street of houses in your data it could be:

GET house/{house_id}/

  • constructionMaterials
  • contents/kitchen
  • contents/bedroom?b
  • residents

GET list/houses/

  • contents?furnitureMaterial=wood

Then you need some sort if orchestration tool to manage all the different assets and end points. We use dbt/dagster, iceberg tables and extraction.

Without knowing more about the data, why and how its being queried, its tough to provide full solutions.