How We Extended Apache DataFusion to Execute One Query Across Many Machines
Datadog, Thursday, October 1st, 2026
Datadog engineers explain Distributed DataFusion, an open-source framework that scales Apache DataFusion queries across machines.
Datadog chose Apache DataFusion as a single composable query engine to replace specialized query systems built for metrics, traces, logs, profiles and security signals.
Because DataFusion runs a query on only one machine, Datadog built Distributed DataFusion, an open-source framework that executes a single query across multiple machines to meet low-latency needs at very large scale.
The post explains the motivation, how the system works and the design decisions behind distributing interactive queries.