10 ms·
A sharded DuckDB on 63 nodes runs 1T row aggregation challenge in 5 sec
- maxmcd 1y agoAre there any open sourced sharded query planners like this? Something that can aggregate queries across many duckdb/sqlite dbs?
- hobofan 1y agoNot directly DuckDB (though I think it might be able to be connected to that), but I think Apache Datafusion Ballista[0] would be a typical modern open source benchmark here. [0]: https://datafusion.apache.org/ballista/contributors-guide/architecture.html https://datafusion.apache.org/ballista/contributors-guide/ar...
- mritchie712 1y agoDeepSeek released smallpond 0 - https://github.com/deepseek-ai/smallpond https://github.com/deepseek-ai/smallpond 1 - https://www.definite.app/blog/smallpond https://www.definite.app/blog/smallpond (overview for data engineers, practical application)
- tgv 1y agoImpressive, but those 63 nodes were "Azure Standard E64pds v6 nodes, each providing 64 vCPUs and 504 GiB of RAM." That's 4000 CPUs and 30TB memory.
- ralegh 1y agoJust noting that 4000 vCPUs usually means 2000 cores, 4000 threads
- RamtinJ95 1y agoAt that scale it cannot be cheaper than just running the same workload on BigQuery or Snowflake or?
- philbe77 1y agoA Standard E64pds v6 costs: $3.744 / hr on demand. At 63 nodes - the cost is: $235.872 / hr - still cheaper than a Snowflake 4XL cluster - costing: 128 credits / hr at $3/credit = $384 / hr.
- ramraj07 1y agoSounds like the equivalent of a 4xl snowflake warehouse, which for such queries would take 30 seconds, with the added benefit of the data being cold stored in s3. Thus you only pay by the minute.
- hobs 1y agoNo, that would be equivalent to 64 4xl snowflake warehouses (though the rest of your point still stands).
- philbe77 1y agoCost-wise, 64 4xl Snowflake clusters would cost: 64 x $384/hr - for a total of: $24,576/hr (I believe)
- __mharrison__ 1y agoWhat was the cost of the duck implementation?
- ramraj07 1y agoApologize for getting it wrong a few orders of magnitude, but thats even more ghastly if its so overpowered and yet takes this long.
- philbe77 1y agoChallenge accepted - I'll try it on a 4XL Snowflake to get actual perf/cost
- shinypenguin 1y agoIs the dataset somewhere accessible? Does anyone know more about the "1T challenge", or is it just the 1B challenge moved up a notch? Would be interesting to see if it would be possible to handle such data on one node, since the servers they are using are quite beefy.
- philbe77 1y agoHi shinypenguin - the dataset and challenge are detailed here: https://github.com/coiled/1trc https://github.com/coiled/1trc The data is in a publicly accessible bucket, but the requester is responsible for any egress fees...
- shinypenguin 1y agoHi, thank you for the link and quick response! :) Do you know if anyone attempted to run this on the least amount of hardware possible with reasonable processing times?
- philbe77 1y agoYes - I also had GizmoSQL (a single-node DuckDB database engine) take the challenge - with very good performance (2 minutes for $0.10 in cloud compute cost): https://gizmodata.com/blog/gizmosql-one-trillion-row-challenge https://gizmodata.com/blog/gizmosql-one-trillion-row-challen...
- simonw 1y agoI suggest linking to that from the article, it is a useful clarification.
- philbe77 1y agoGood point - I'll update it...
- achabotl 1y agoThe One Trillion Row Challenge was proposed by Coiled in 2024. https://docs.coiled.io/blog/1trc.html https://docs.coiled.io/blog/1trc.html
- MobiusHorizons 1y ago> Once trusted, each worker executes its local query through DuckDB and streams intermediate Arrow IPC datasets back to the server over secure WebSockets. The server merges and aggregates all results in parallel to produce the final SQL result—often in seconds. Can someone explain why you would use websockets in an application where neither end is a browser? Why not just use regular sockets and cut the overhead of the http layer? Is there a real benefit I’m missing?
- philbe77 1y agoHi MobiusHorizons, I happened to use websockets b/c it was the technology I was familiar with. I will try to learn more about normal sockets to see if I could perhaps make them work with the app. Thanks for the suggestion...
- DanielHB 1y agoif you really want maximum performance maybe consider using CoAP for node-communication: https://en.wikipedia.org/wiki/Constrained_Application_Protocol https://en.wikipedia.org/wiki/Constrained_Application_Protoc... It is UDP-based but adds handshakes and retransmissions. But I am guessing for your benchmark transmission overhead isn't a major concern. Websockets are not that bad, only the initial connection is HTTP. As long as you don't create a ton of connections all the time it shouldn't be much slower than a TCP-based socket (purely theoretical assumption on my part, I never tested).
- gopalv 1y ago> will try to learn more about normal sockets to see if I could perhaps make them work with the app. There's a whole skit in the vein of "What have the Romans ever done for us?" about ZeroMQ[1] which has probably lost to the search index now. As someone who has held a socket wrench before, fought tcp_cork and dsack, Websockets isn't a bad abstraction to be on top of, especially if you are intending to throw TLS in there anyway. Low level sockets is like assembly, you can use it but it is a whole box of complexity (you might use it completely raw sometimes like a tickle ack in the ctdb[2] implementation). [1] - https://news.ycombinator.com/item?id=32242238 https://news.ycombinator.com/item?id=32242238 [2] - https://linux.die.net/man/1/ctdb https://linux.die.net/man/1/ctdb
- nodesocket 1y ago> Each GizmoEdge worker pod was provisioned with 3.8 vCPUs (3800 m) and 30 GiB RAM, allowing roughly 16 workers per node—meaning the test required about 63 nodes in total. How was this node setup chosen? Specially 3.8 vCPU and 30 GiB RAM per? Why not just run 16 workers total using the entire 64 vCPU and 504 GiB of memory each?
- philbe77 1y agoHi nodesocket - I tried to do 4 CPUs per node, but Kubernetes takes a small (about 200m) CPU request amount for daemon processes - so if you try to request 4 (4000m) CPUs x 16 - you'll spill one pod over - fitting only 15 per node. I was out of quota in Azure - so I had to fit in the 63 nodes... :)
- nodesocket 1y agoBut why split up a vm into so many workers instead of utilizing the entire vm as a dedicated single worker? What’s the performance gain and strategy?
- philbe77 1y agoI'm not exactly sure yet. My goal was to not have the shards be too large so as to be un-manageable. In theory - I could just have had 63 (or 64) huge shards - and 1 worker per K8s node, but I haven't tried it. There are so many variables to try - it is a little overwhelming...
- nodesocket 1y agoWould be interesting to test. I’m thinking there may not be a benefit to having so many workers on a vm instead of just the entire vm resources as a single worker. Could be wrong, but that would be a bit surprising.
- boshomi 1y ago>“In our talk, we will describe the design rationale of the DuckLake format and its principles of simplicity, scalability, and speed. We will show the DuckDB implementation of DuckLake in action and discuss the implications for data architecture in general. Prof. Hannes Mühleisen, cofounder of DuckDB: [DuckLake - The SQL-Powered Lakehouse Format for the Rest of Us by Prof. Hannes Mühleisen](https://www.youtube.com/watch?v=YQEUkFWa69o https://www.youtube.com/watch?v=YQEUkFWa69o) (53 min) Talk from Systems Distributed '25: https://systemsdistributed.com https://systemsdistributed.com
- ferguess_k 1y agoWait until you see a 800-line Tableau query that joins TB data with TB data /s
- NorwegianDude 1y agoThis is very silly. You're not doing the challenge if you do the work up front. The idea is that you start with a file and the goal is to get the result as fast as possible. How long did it take to distribute and import the data to all workers, what is the total time from file to result? I can do this a million times faster on one machine, it just depends on what work I do up front.
- philbe77 1y agoYou should do it then, and post it here. I did do it with one machine as well: https://gizmodata.com/blog/gizmosql-one-trillion-row-challenge https://gizmodata.com/blog/gizmosql-one-trillion-row-challen...
- NorwegianDude 1y agoNobody cares if I can do it a million times faster, everyone can. It's cheating. The whole reason you have to account for the time you spend setting it up is so that all work spent processing the data is timed. Otherwise we can just precomputed the answer and print it on demand, that is very fast and easy. Just getting it into memory is a large bottleneck in the actual challenge. If I first put it into a DB with statistics that tracks the needed min/max/mean then it's basically instant to retrieve, but also slower to set up because that work needs to be done somewhere. That's why the challenge is time from file to result.
- ta12653421 1y agoWhen reading such extreme numbers, I'm always thinking what I may be doing wrong, when my MSSQL based CRUD application warms up its caches with around 600.000 rows and it takes 30 seconds to load them from DB into RAM on my 4x3GHz machine :-D Maybe I'm missing something fundamental here
- dgan 1y agoI also had misfortune working with MSSQL is it was so so unbearably slow, because i couldnt upload data in bulk. I guess its forbidden technology
- Foobar8568 1y agoOr you didn't use MSSQL properly, there are at least 2 or 3 ways to do bulk upload on MS SQL, not sure in today era.
- dgan 1y agoMaybe? Don't know. I never had problemes bulk uploading into Postgres tho, it's right there in documentation and I don't have to have a weird executable on my corporately castrated laptop
- Foobar8568 1y agohttps://learn.microsoft.com/en-us/sql/t-sql/statements/bulk-insert-transact-sql?view=sql-server-ver17 https://learn.microsoft.com/en-us/sql/t-sql/statements/bulk-... That's one way, another was BCP. But yeah if you are using python and loading row by row, or a large amount into a large table that has a clustered index, chances are that it'll be dead slow but that's expected.
- dgan 1y agoThe documentation you providee requires for the file to be present on the server side, not the client side, which is very dufferent from postgres. As for BCP executable, i couldnt find a way for it to accept any type of date[time] at all
- mosselman 1y agoAre there any good instructions somewhere on how to set this up? As in not 63 nodes. But a distributed duckdb instance
- philbe77 1y agoHi mosselman, GizmoEdge is not open-source. DeepSeek has "smallpond" however, which is open-source: https://github.com/deepseek-ai/smallpond https://github.com/deepseek-ai/smallpond I plan on getting GizmoEdge to production-grade quality eventually so folks can use it as a service or licensed software. There is a lot of work to do, though :)
- djhworld 1y agoInteresting and fun > Workers download, decompress, and materialize their shards into DuckDB databases built from Parquet files. I'm interested to know whether the 5s query time includes this materialization step of downloading the files etc, or is this result from workers that have been "pre-warmed". Also is the data in DuckDB in memory or on disk?
- philbe77 1y agohi djhworld. The 5s does not include the download/materialization step. That parts takes the worker about 1 to 2 minutes for this data set. I didn't know that this was going on HackerNews or would be this popular - I will try to get more solid stats on that part, and update the blog accordingly. You can have GizmoEdge reference cloud (remote) data as well, but of course that would be slower than what I did for the challenge here... The data is on disk - on locally mounted NVMe on each worker - in the form of a DuckDB database file (once the worker has converted it from parquet). I originally kept the data in parquet, but the duckdb format was about 10 to 15% faster - and since I was trying to squeeze every drop of performance - I went ahead and did that... Thanks for the questions. GizmoEdge is not production yet - this was just to demonstrate the art of the possible. I wanted to divide-and-conquer a huge dataset with a lot of power...
- philbe77 1y agoI've since learned (from a DuckDB blog) - that DuckDB seems to do better when the XFS filesytem. I used ext4 for this, so I may be able to get another 10 to 15% (maybe!). DuckDB blog: https://duckdb.org/2025/10/09/benchmark-results-14-lts https://duckdb.org/2025/10/09/benchmark-results-14-lts
- sammy2255 1y agoHow would a 63 node Clickhouse cluster compare? >:)
- lolive 1y agoWhy doesn't such large-scale test the big feature everyone needs, which is inner join at scale?
- philbe77 1y agoThis is something we are trying to take a novel approach to as well. We have a video demonstrating some TPC-H SF10TB queries which perform inner joins, etc. - with GizmoEdge as well: https://www.youtube.com/watch?v=hlSx0E2jGMU https://www.youtube.com/watch?v=hlSx0E2jGMU
- lolive 1y agoDoes that study go into the global vision of DuckLake ?
- 1a527dd5 1y agoThe title buries the lede a little > Our cluster ran on Azure Standard E64pds v6 nodes, each providing 64 vCPUs and 504 GiB of RAM. Yes, I would _expect_ when each node has that kind of power it should return very impressive speeds.
- vysakh0 1y agoDuckdb is an excellent OLAP db, I have had customers who had s3 data lake of parquet and use databricks or other expensive tool, when they could easily use duckdb.. Given we have cursor/claude code, it is not that hard for lot of use cases, I think the lack of documentation on how duckdb functions -- in terms of how it loads these files etc are some of the reasons companies are not even trying to adopt duckdb. I think blogs like this is a great testament for duckdb's performance!
- mrtimo 1y agoI have experience with duckDB but not databricks... from the perspective of a company, is a tool like databricks more "secure" than duckdb? If my company adopts duckdb as a datalake, how do we secure it?
- rapatel0 1y agoDuckdb can run as a local instance that points to parquet files in a n s3 bucket. So your "auth" can live on the layer that gives permissions to access that bucket.
- lopatin 1y agoDuckDB is great but it’s barely OLAP right? A key part of OLAP is “online”. Since the writer process blocks any other processes from doing reads, calling it OLAP is a stretch I think.
- ansgri 1y agoIsn't the Online part here about getting results immediately after query, as opposed to overnight batch reports? So if you don't completely overwhelm DuckDB with writes, it still qualifies. The quality you're describing is something like "realtime analytics", and is a whole another category: Clickhouse doesn't qualify (batching updates, merging etc. — but it's clearly OLAP), Druid does.
- sdairs 1y ago
- sdairs 1y agoPretty big caveat; 5 seconds AFTER all data has been loaded into memory - over 2 minutes if you also factor reading the files from S3 and loading memory. So to get this performance you will need to run hot: 4000 CPUs and and 30TB of memory going 24/7.
- CaptainOfCoit 1y agoYeah, pretty misleading it feels like. For background, here is the initial ideation of the "One Trillion Row Challenge" challenge this submission originally aimed to participate in: https://docs.coiled.io/blog/1trc.html https://docs.coiled.io/blog/1trc.html
- lumost 1y agoIt does make me wonder whether all of the investment in hot-loading of GPU infrastructure for LLM workloads is portable to databases. 30TB of GPU memory will be roughly 200 B200 cards or roughly 1200 per hour compared to the $240/hour pricing for the CPU based cluster. The GPU cluster would assuredly crush the CPU cluster with a suitable DB given it has 80x the FP32 FLOP capacity. You'd expect the in-memory GPU solution to be cheaper (assuming optimized software) with a 5x growth in GPU memory per card, or today if the workload can be bin-packed efficiently.
- eulgro 1y agoDo databases do matrix multiplication? Why would they even use floats?
- lumost 1y agolot's of columns are float valued, GPU tensor cores can be programmed to do many operations between different float/int valued vectors. Strings can also be processed in this manner as they are simply vectors of integers. NVidia publishes official TPC benchmarks for each GPU release. The idea of a GPU database has been reasonably well explored, they are extremely fast - but have been cost ineffective due to GPU costs. When the dataset is larger than GPU memory, you also incur slowdowns due to cycling between CPU and GPU memory.
- afpx 1y agoI’ve never used DuckDB, but I was surprised by the 30 GiB of memory. Many years ago when I used to use EMR a lot, I would go for > 10 TiB of RAM to keep all the data in memory and only spill over to SSD on big joins.
- kwillets 1y agoThis is fun, but I'm confused by the architecture. Duckdb is based on one-off queries that can scale momentarily and then disappear, but this seems to run on k8s and maintain a persistent distributed worker pool. This pool lacks many of the features of a distributed cluster such as recovery, quorum, and storage state management, and queries run through a single server. What happens when a node goes down? Does it give up, replan, or just hang? How does it divide up resources between multiple requests? Can it distribute joins and other intermediate operators? I have a soft spot in my heart for duckdb, but its uniqueness is in avoiding the large-scale clustering that other engines already do reasonably well.
- fHr 1y agoBetter and worth more then all the quantum bs I have to listen to.
- hbarka 1y agoSELECT COUNT(DISTINCT) has entered the challenge.
- philbe77 1y agogood point :) - we can re-aggregate HyperLogLog (HLL) sketches to get a pretty accurate NDV (Count Distinct) - see Query.farm's DataSketches DuckDB extension here: https://github.com/Query-farm/datasketches https://github.com/Query-farm/datasketches We also have Bitmap aggregation capabilities for exact count distinct - something I worked with Oracle, Snowflake, Databricks, and DuckDB labs on implementing. It isn't as fast as HLL - but it is 100% accurate...
- fifilura 1y agoI remember BigQuery had Distinct with HLL accuracy 10 years ago but rather quickly replaced it with actual accuracy. How would you compare this solution to BigQuery?
- deleted 1y ago[deleted]
- up2isomorphism 1y agoSensational title, a reflection of “attention is all you need”.(pun intended)
- fifilura 1y agoIsn't Trino built for exactly this, without the quirky workarounds?
- peter_d_sherman 1y ago>"The GizmoEdge Server receives a SQL query from the client, parses it, and generates two statements: o A worker SQL to execute on each distributed node o A combinatorial SQL to run server-side for final aggregation" A specific instance of MapReduce (using SQL!): https://en.wikipedia.org/wiki/MapReduce https://en.wikipedia.org/wiki/MapReduce