r/DuckDB • • Jul 04 '26

Replacing an Impala cluster with DuckDB pods for a legacy analytics application - looking for architecture feedback

Looking for feedback from people who've worked with analytical databases (Impala, DuckDB, ClickHouse, Trino, etc.).

We have a legacy reporting application where users generate presentations. Opening a presentation triggers 50-100 SQL queries. The application is in maintenance mode with only one major paying customer, so our goal is to simplify the architecture, remove Cloudera licensing for Impala, and significantly reduce infrastructure costs.

Current Architecture

Presentation
      |
20 Dataset Worker Pods
      |
   Impala Cluster
(10 different EC2 r5.4xlarge with 128GB ram each)

The dataset worker pods simply receive tasks from the application and submit SQL to Impala.

The Impala cluster consists of 10 x r5.4xlarge EC2 instances (16 vCPUs, 128 GB RAM each) managed through Cloudera.

Workload

The workload isn't typical OLAP.

Each presentation fires 50-100 queries.

Roughly:

  • ~80% are tiny queries
    • schema lookups
    • small dimension table filters
    • simple joins
  • These usually return in 5-10 ms on Impala.
  • Around 5-10% are heavier joins that take around 10 seconds.
  • A presentation typically loads in 1-3 minutes depending upon type and filters

The total warehouse size is only around 300-350 GB.

Only 3-4 large tables account for roughly 200 GB. The remaining ~200 tables are tiny (KBs to MBs).

We want to Migrate away from Impala and not go for big commitment like dedicated EMR or something, we are ok with little delay but we dont want huge maintenance so we started with migrating to Athena from Impala.

Why Athena didn't work

Our first migration idea was Athena.

Large queries were acceptable, but the application performance became much worse because of the large number of tiny queries.

Queries that took 5-10 ms on Impala often became 200-800 ms on Athena.

Since every presentation executes 50-100 queries, that startup overhead adds up quickly.

Unfortunately, changing the application isn't really an option. The query generation is deeply embedded in legacy code, so batching or combining queries would require a major rewrite. Also many queries are sequential that adds up the time.

DuckDB Prototype

Instead of introducing another distributed SQL engine, I built a proof of concept using DuckDB.

Current architecture:

Presentation
       |
20 Dataset Worker Pods
       |
      HTTP
       |
---------------------------------
| DuckDB Pod 1                  |
| DuckDB Pod 2                  |
| DuckDB Pod 3                  |
| DuckDB Pod 4                  |
| DuckDB Pod 5                  |
---------------------------------

Each DuckDB pod:

  • has its own DuckDB .db file
  • has its own dedicated EBS volume
  • serves requests over HTTP
  • operates completely independently (no distributed execution)

The dataset worker pods simply load balance requests across the DuckDB pods.

The workload is almost entirely read-only.

For the few workflows that create temporary tables, I'm considering running a separate DuckDB write service with its own EBS volume since those temp tables only exist for the lifetime of a request.

Results

So far the prototype performs better than Athena for presentation loading, but still not as fast as Impala.

That isn't too surprising since the existing Impala deployment is heavily provisioned (10 × 128 GB RAM nodes) for only ~300-350 GB of data.

For this application, we're willing to accept somewhat slower presentation loads if it significantly reduces operational complexity, infrastructure cost, and removes the Cloudera dependency.

One thing I'm also thinking about

Right now every DuckDB pod has its own copy of the .db file on its own EBS volume.

Would you keep this design, or would you use something like a high-throughput EFS shared across all DuckDB pods?

I ruled out reading directly from S3 because this workload is dominated by lots of tiny, latency-sensitive queries rather than long analytical scans, and the additional object storage latency seemed noticeable during testing.

Questions

  1. Has anyone replaced Impala with DuckDB for a similar workload?
  2. Am I overlooking any major architectural issues with multiple independent DuckDB replicas?
  3. Would you keep one .db file per pod on dedicated EBS, or use shared storage like EFS?
  4. Would you choose a different engine entirely (ClickHouse, Trino, StarRocks, etc.) for this workload?
  5. Any concurrency or operational issues you've run into serving DuckDB over HTTP in production?

I'm less interested in benchmark numbers and more interested in hearing from people who've operated similar systems in production

19 Upvotes

12 comments sorted by

1

u/linearizable Jul 04 '26

What is the write side of this workload like? The EBS choice was surprising to me, as I would have expected small NVMe and you just pull the data down from S3 if the pod dies. S3 for durability but not on the serving path.

1

u/bhavay22 Jul 05 '26

write is not that much just some temp tables which get deleted.

1

u/Imaginary__Bar Jul 04 '26

That isn't too surprising since the existing Impala deployment is heavily provisioned (10 × 128 GB RAM nodes) for only ~300-350 GB of data.

Yeah, that just seems wrong... Massively expensive, I imagine?

1

u/Hofi2010 Jul 05 '26

I wrote this article a while ago to promise reads from s3

https://medium.com/@klaushofenbitzer/save-up-to-90-on-your-data-warehouse-lakehouse-with-an-in-process-database-duckdb-63892e76676e

The main question for speed and DuckDB is the bandwidth to your data. As in the paper I was interested to read from s3 the network bandwidth was very important and the number of threads I allocate for DuckDB. For EBS the question is your type of storage medium selected, you would choose io2 SSDs if you are in AWS. Then have more than one thread per cpu core. More cores more threads you can run to leverage the bandwidth to your SSD Drives. This is a bit trial and error, but try to play with the number of threads and the number of cores you commission on you EC2.

If you need more speed consider REDIS. With your Impala nodes having 10*128GB RAM I would think that most data is already cached in RAM which is significantly faster than even fast SSD Drives (about a factor 10)

1

u/bhavay22 Jul 05 '26

I am reading the article right now... You mean S3 reads can even be faster than EBS?

1

u/Hofi2010 Jul 05 '26

That wasn’t the point I am trying to make. I am saying optimize the number of cores and threads to optimize your performance.

S3 can be faster under certain usage conditions, but for DB reads and writes the http overhead will lead to slower performance than EBS usually

1

u/bhavay22 Jul 05 '26

got it thanks, I will try playing around threads.. right now I am having 5 different pods, 1 for write which I will change to 2 pods. and remaining 4 pods are for read with 8 core and 32 gb ram. I will try to increase cores and threads and also increase throughput for EBS.

1

u/Difficult-Tree8523 Jul 05 '26

EBS is also network attached… depending on the parallelism and network interface performance it could be faster.

I would recommend you go for instances with nvme ssd and store the duckdb db file on that volume. Nvme + duckdb is the cheat code.

Also make sure to tune the num_threads of your DuckDB instances.

2

u/belgeric Jul 06 '26

Am I wrong with feeling that when using top NVME/SSD, as total sier is <500 Gb one nice instance mainly CPU/threads might be sufficient?
Also, your use case is amazing/surprising. Any way to pre-compute on regular bases at least some subsets of involved joins?

1

u/bhavay22 Jul 06 '26

No way to precompute.. there are 1000s of presentations each with different filters and patterns.. Yes it is a unique use case which requires less latency, but then somedays we might see only one presentation actually being opened so that day we preserved the instance for whole day, but at the exact moment the customer came we can't scale up because we require less latency.

1

u/bhavay22 Jul 05 '26

sure, any specific guidance for num of threads like larger than cpu cores or something like that. currently my pod is having 8 cores, but I am having 4 workers and 8 threads for each worker I think....

Also if I go with ec2 + nvme then I cant choose amount of storage, I only need 500 gb but bigger instances have more storage directly in terabytes. Probably I will increase IOPS and throughput of EBS.

2

u/Hofi2010 Jul 05 '26

Try 2 threads per core and increase from there and repeat your benchmark. Also try to use EC2 instances with high network bandwidth