Skip to content

Repository files navigation

Duckherder - DuckDB Remote and Distributed Execution Extension

Duckherder is a DuckDB extension built upon storage extension that enables (certain) distributed query execution across multiple worker nodes using Arrow Flight for efficient data transfer. It allows you to seamlessly work with remote tables and execute queries in parallel across distributed workers while maintaining DuckDB's familiar SQL interface.

Warning

This is a personal project and has nothing to do with DuckLabs. It is not affiliated with, endorsed by, or associated with DuckDB Labs in any way.

WIP Disclaimer

This repository is currently a work in progress, and not yet for production usage. Feel free to play around with it, give me feedback, and ping me for feature request and collaboration!

Table of Contents

Overview

Duckherder implements a client-server architecture with distributed query execution:

  • Client (Duckherder): Coordinates queries, manages remote table references, and initiates distributed execution
  • Server: Consists of two components:
    • Driver Node: Analyzes queries, creates partition plans, and coordinates worker execution
    • Worker Nodes: Execute partitioned queries on local data and return results via Arrow Flight
  • Distributed Executor: Runs on the driver node, partitions queries based on DuckDB's physical plan analysis, and distributes tasks to workers

The extension transparently handles query routing, allowing you to run CREATE, SELECT, INSERT, DELETE, and ALTER operations on remote tables through DuckDB's storage extension interface.

Architecture

High-Level System Architecture

┌─────────────────────────────────────────┐
│         Client DuckDB Instance          │
│  ┌───────────────────────────────────┐  │
│  │   Duckherder Catalog (dh)         │  │
│  │   - Remote table references       │  │
│  │   - Query routing logic           │  │
│  └───────────────────────────────────┘  │
│              │                          │
│              │ Arrow Flight Protocol    │
│              │                          │
└──────────────┼──────────────────────────┘
               │
┌──────────────┴──────────────────────────┐
│         Server DuckDB Instance(s)       │
│  ┌───────────────────────────────────┐  │
│  │   Driver Node                     │  │
│  │   - Distributed executor          │  │
│  │   - Query plan analysis           │  │
│  │   - Task coordination             │  │
│  └───────────┬───────────────────────┘  │
│              │                          │
│  ┌───────────┴───────────────────────┐  │
│  │   Worker Nodes                    │  │
│  │   - Actual table data             │  │
│  │   - Partitioned execution         │  │
│  │   - Result streaming              │  │
│  └───────────────────────────────────┘  │
└─────────────────────────────────────────┘

Object-Storage-Backed Deployment Model

The target architecture for integrating distributed execution with an object-storage-backed DuckDB database follows a single-writer, multiple-reader model:

  • The Driver/Coordinator owns the only write authority. In the initial design, the Writer runs inside the Driver process. It owns the read-write database handle and writer lease, applies DDL/DML, and publishes new checkpoints or snapshots to object storage.
  • All worker nodes are readers. They execute partitioned read tasks against the snapshot selected by the coordinator and never publish database changes.
  • The client is a SQL and session endpoint. It registers storage and cluster configuration, attaches the database, and submits queries without coordinating individual workers or writing database files directly.
  • Object storage is the source of truth for database data, manifests, and committed snapshots.
+------------------+    SQL / transaction     +------------------------------------+
| Client CLI / App | -----------------------> | Driver / Coordinator               |
| - ATTACH         | <----------------------- |                                    |
| - Session        |         results          |  +------------------------------+  |
+------------------+                          |  | Gateway + Result Merger      |  |
                                              |  +------------------------------+  |
                                              |  | Catalog + Query Planner      |  |
                                              |  | table ID + schema + snapshot |  |
                                              |  +------------------------------+  |
                                              |  | Transaction Coordinator      |  |
                                              |  | BEGIN / COMMIT / ROLLBACK    |  |
                                              |  +------------------------------+  |
                                              |  | Writer (READ-WRITE)          |  |
                                              |  | writer lease + epoch         |  |
                                              |  | DDL/DML + checkpoint/publish |  |
                                              |  +------------------------------+  |
                                              +----------+-------------+-----------+
                                                         |             |
                                read tasks + snapshot ID |             | write +
                                                         |             | publish
                                                         v             v
                                              +----------------+  +------------------+
                                              | Worker Pool    |  | Object Storage   |
                                              | READ-ONLY      |  | data + manifests |
                                              | Worker 1 ... N |->| + checkpoints    |
                                              +-------+--------+  +------------------+
                                                      |             ^
                                                      |             |
                                                      +-------------+
                                                   read pinned snapshot

                                              partial results return to
                                              Gateway + Result Merger

Each query is pinned to one committed snapshot before tasks are sent to workers. A newly committed write is visible to new queries, while already running queries continue reading their original snapshot. The coordinator remains the commit authority even if the Writer is later moved from the Driver process into a dedicated writer process.

This section describes the intended object-storage-backed distributed architecture. Transaction coordination, snapshot propagation, and direct worker reads from shared object storage are still work in progress.

Distributed Execution Flow

┌─────────────┐
│ COORDINATOR │ Extract query plan
└──────┬──────┘
       │
       ├─ Phase 1: Extract & validate logical plan
       ├─ Phase 2: Analyze DuckDB's natural parallelism
       ├─ Phase 3: Create partition plans
       ├─ Phase 4: Prepare result schema
       │
       ├──────────┬──────────┬──────────┐
       ▼          ▼          ▼          ▼
   ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐
   │WORKER 0│ │WORKER 1│ │WORKER 2│ │WORKER N│
   └───┬────┘ └───┬────┘ └───┬────┘ └───┬────┘
       │          │          │          │
       │ Execute  │ Execute  │ Execute  │ Execute
       │ partition│ partition│ partition│ partition
       │ (Local)  │ (Local)  │ (Local)  │ (Local)
       │          │          │          │
       └──────────┴──────────┴──────────┘
              │
              ▼
       ┌─────────────┐
       │ COORDINATOR │  Merge results
       │   Combine   │  Smart aggregation
       │  Finalize   │  Return final result
       └─────────────┘

Distributed Execution System

Partitioning Strategy

The distributed executor analyzes DuckDB's physical plan to create optimal partitions:

1. Natural Parallelism Analysis

  • Queries EstimatedThreadCount() from DuckDB's physical plan
  • Understands how many parallel tasks DuckDB would naturally create
  • Extracts cardinality and operator information

2. Partition Strategy Selection (Priority Order)

The system uses a three-tier hierarchy to select the optimal partitioning strategy:

Priority Condition Strategy Example Predicate Notes
1st Row groups detected (cardinality known) ROW GROUP-BASED WHERE rowid BETWEEN 0 AND 245759 Assigns whole row groups to tasks
2nd TABLE_SCAN + ≥100 rows/worker RANGE-BASED WHERE rowid BETWEEN 0 AND 2499 Contiguous ranges, not row-group-aligned
3rd Small table / Non-TABLE_SCAN / Unknown cardinality MODULO WHERE (rowid % node id) = 0 Fallback for all other cases. TODO(hjiang): for small table no need to distribute, should directly execute on driver node.

Row group-based partitioning (preferred):

  • Aligned with DuckDB's storage structure (122,880 rows per group)
  • Tasks process complete row groups (e.g., Task 0: groups [0,1], Task 1: groups [2,3])
  • Optimal cache locality and I/O efficiency
  • Respects DuckDB's natural parallelism boundaries

Range-based partitioning (fallback):

  • Used when row groups can't be detected but cardinality is known
  • Contiguous row access
  • Not aligned with row group boundaries

Modulo partitioning (final fallback):

  • Scattered row access (poor cache locality)
  • Works for any table size and operator type
  • Simple and always correct

3. How Row Group-Based Partitioning Works

When row groups are detected (the common case), the system:

  • Calculates total row groups: cardinality / 122,880
  • Assigns whole row groups to tasks (e.g., Task 0 gets groups [0,1], Task 1 gets groups [2,3])
  • Converts row group ranges to rowid ranges: rowid BETWEEN (rg_start * 122880) AND (rg_end * 122880 - 1)
  • Workers scan complete row groups for optimal performance

DuckDB Thread Model → Distributed Model Mapping

DuckDB Thread Model Distributed Model Implementation
Thread Worker Node Physical machine/process
LocalSinkState Worker Result QueryResult → Arrow batches
Sink() Worker Execute HandleExecutePartition()
GlobalSinkState ColumnDataCollection Coordinator's result collection
Combine() CollectAndMergeResults() Merging worker outputs
Finalize() Return MaterializedQueryResult Final result to client

Installation

Building from Source

Prerequisites

  1. Install VCPKG for dependency management:
git clone https://github.com/Microsoft/vcpkg.git
./vcpkg/bootstrap-vcpkg.sh
export VCPKG_TOOLCHAIN_PATH=`pwd`/vcpkg/scripts/buildsystems/vcpkg.cmake
  1. Build the extension:
# Clone the repo.
git clone --recurse-submodules https://github.com/dentiny/duckdb-distributed-execution.git
cd duckdb-distributed-execution

# Build with release mode.
export VCPKG_TOOLCHAIN_PATH=/path/to/vcpkg/scripts/buildsystems/vcpkg.cmake
CMAKE_BUILD_PARALLEL_LEVEL=$(nproc) make

The build produces:

  • ./build/release/duckdb - DuckDB shell with extension pre-loaded
  • ./build/release/test/unittest - Test runner
  • ./build/release/extension/duckherder/duckherder.duckdb_extension - Loadable extension
  • ./build/release/extension/duckdb_object_storage/duckdb_object_storage.duckdb_extension - Object-storage filesystem extension
  • ./build/release/distributed_server - Standalone distributed driver node
  • ./build/release/distributed_worker - Standalone distributed worker node

Running Standalone Server and Workers

For production deployments or testing distributed execution across multiple machines, you can run standalone server and worker processes:

Starting the Driver/Coordinator Server

# Start the driver server on default host (0.0.0.0) and port (8815) without workers
./build/release/distributed_server

# Start the driver server on a specific host and port
./build/release/distributed_server 192.168.1.100 8815

# Start the driver server with 4 local workers (for single-machine distributed execution)
./build/release/distributed_server 0.0.0.0 8815 4

Starting Standalone Worker Nodes

For multi-machine setups, start worker nodes on separate machines:

# Start a worker on default host (0.0.0.0) and port (8816) with default worker ID (worker-1)
./build/release/distributed_worker

# Start a worker on a specific port with custom worker ID
./build/release/distributed_worker 0.0.0.0 8817 worker-2

Object Storage

The build can attach a native DuckDB database stored in S3-compatible object storage. The currently supported deployment is one writer with one or more read-only processes. In the target distributed deployment, the Driver/Coordinator owns that writer and every worker is a read-only process. See the S3 single-writer/read-only-reader test and local object-storage test for executable examples.

For local S3-compatible testing, the RustFS helper starts a persistent Docker-backed service, creates the test bucket, and can run an extension write/read smoke test:

./scripts/local-rustfs.sh start   # S3 API :19000, console :19001
./scripts/local-rustfs.sh test    # write, checkpoint, reattach read-only, query
./scripts/local-rustfs.sh stop    # keeps the Docker volume
./scripts/local-rustfs.sh reset   # deletes the test data

Accessing Local RustFS

The helper uses these defaults:

  • S3 API endpoint: http://127.0.0.1:19000
  • Web console: http://127.0.0.1:19001
  • Bucket: duckherder
  • Root prefix: duckherder-test
  • Access key: rustfsadmin
  • Secret key: rustfsadmin
  • Region: us-east-1
  • URL style: path-style

Open the web console and sign in with the access and secret keys above, or use any S3-compatible client. For example, with the AWS CLI:

AWS_ACCESS_KEY_ID=rustfsadmin \
AWS_SECRET_ACCESS_KEY=rustfsadmin \
AWS_DEFAULT_REGION=us-east-1 \
aws --endpoint-url http://127.0.0.1:19000 s3 ls s3://duckherder/duckherder-test/

Run ./scripts/local-rustfs.sh status at any time to print the active endpoint, bucket, root, credentials, and matching DuckDB secret SQL. All settings can be overridden with RUSTFS_* environment variables. The same overrides must be used for subsequent start, status, test, and stop commands.

Configuring the DuckDB S3 Secret

Load an extension that registers DuckDB's S3 secret type, then create a config-provider secret whose scope covers the bucket and root prefix:

LOAD cache_httpfs;

CREATE OR REPLACE SECRET local_rustfs (
    TYPE S3,
    PROVIDER CONFIG,
    KEY_ID 'rustfsadmin',
    SECRET 'rustfsadmin',
    REGION 'us-east-1',
    ENDPOINT '127.0.0.1:19000',
    USE_SSL false,
    URL_STYLE 'path',
    SCOPE 's3://duckherder/duckherder-test'
);

The ENDPOINT must not include http://; USE_SSL false selects HTTP. URL_STYLE 'path' is required for this local endpoint. The DATA_PATH must fall within the secret's SCOPE.

To use this storage through Duckherder, start a control node and attach a named database:

SELECT duckherder_start_local_server(8815);

ATTACH 'localhost:8815/database.db' AS dh (
    TYPE duckherder,
    DATA_PATH 's3://duckherder/duckherder-test/manual',
    SECRET 'local_rustfs'
);

Duckherder resolves the secret and sends the required S3 configuration to the control node and workers, which create temporary S3 secrets before initializing duckdb_object_storage. The current Flight transport does not encrypt these credentials; use S3 attachments only on a trusted network until TLS transport is supported.

Usage

Local Server Management

-- Start a distributed server on port 8815.
SELECT duckherder_start_local_server(8815);

-- Stop the local distributed server.
SELECT duckherder_stop_local_server();

Attach to the Server

-- READ_WRITE is the default.
ATTACH DATABASE 'localhost:8815' AS dh (TYPE duckherder);

-- READ_ONLY must be explicit.
ATTACH DATABASE 'localhost:8815' AS dh (TYPE duckherder, READ_ONLY);

The ATTACH path does not store Duckherder data. Duckherder always backs its local DuckCatalog metadata cache with an in-memory database; the control node remains authoritative and no local catalog or table file is persisted. During ATTACH, Duckherder discovers remote enum types and tables and loads their definitions into this cache. Tables created through the attachment are created remotely first and added to the cache automatically. The path must be a remote host:port endpoint (an explicit grpc://host:port URL is also accepted); :memory: and filesystem paths are not supported by the Duckherder storage type.

Connection and Access Model

A Duckherder attachment is visible to every DuckDB connection that shares the same DatabaseInstance. Each connection that is allowed to use the attachment owns a separate Flight client, control-node registration, and server-side DuckDB connection.

For a READ_WRITE attachment, exactly one DuckDB connection owns the attachment at a time. That connection may read and write; every operation from another connection is rejected:

-- Connection 1: creates and owns the READ_WRITE attachment.
ATTACH DATABASE 'localhost:8815' AS dh (TYPE duckherder);
SELECT * FROM dh.my_table; -- allowed
INSERT INTO dh.my_table VALUES (1); -- allowed

-- Connection 2, while Connection 1 is alive:
SELECT * FROM dh.my_table;
-- Error: the READ_WRITE attachment belongs to another connection.

When the owner connection closes, its Flight client and server registration are closed. The next connection that uses the attachment becomes its new owner.

For a READ_ONLY attachment, every DuckDB connection may read through its own independent client, while mutations are rejected:

-- Connection 1: creates a READ_ONLY attachment.
ATTACH DATABASE 'localhost:8815' AS dh (TYPE duckherder, READ_ONLY);

-- Connection 1 and Connection 2: both allowed, using different Flight clients.
SELECT * FROM duckherder_get_query_execution_stats();

-- Either connection:
INSERT INTO dh.my_table VALUES (1);
-- Error: Duckherder client is read-only.

The control node admits at most one writable client registration and any number of read-only registrations. DETACH closes every client belonging to the attachment and releases the writable slot. Attachments renew a short control-node lease, so a crashed writer is reclaimed after its lease expires. Client and control node must use the same protocol version because role registration is not compatible with older binaries.

Remote Table Registration

ATTACH discovery and CREATE TABLE register remote tables automatically. Duckherder does not expose a manual registration API. A mapping can still be removed explicitly:

-- Unregister a remote table mapping.
-- Syntax: duckherder_unregister_remote_table(local_table_name)
PRAGMA duckherder_unregister_remote_table('my_table');

Load Extensions on Server

-- Load an extension on the remote server
SELECT duckherder_load_extension('parquet');
SELECT duckherder_load_extension('json');

Driver/Worker Node Registration

Starting and Registering Driver Nodes

For distributed execution with external coordination, you can register a driver node location. Unlike worker nodes which can be added incrementally, only one driver node can be registered at a time - registering a new driver will automatically replace any previously registered driver.

Option 1: Start local driver server (for testing, debugging and development)

-- Start a local server on default host (0.0.0.0) and port (8815)
SELECT duckherder_start_local_server(8815);

Option 2: Register external driver (for multi-machine setups)

-- Register an external driver node running on a different machine
-- Syntax: duckherder_register_or_replace_driver(driver_id, location)
SELECT duckherder_register_or_replace_driver('driver-main', 'grpc://192.168.1.100:8815');

-- Replacing an existing driver node (automatic replacement)
SELECT duckherder_register_or_replace_driver('driver-backup', 'grpc://192.168.1.200:8815');

Option 3: Using the standalone distributed server executable

# Start the driver server on default host (0.0.0.0) and port (8815)
./build/release/distributed_server

# Start the driver server on a specific host and port
./build/release/distributed_server 192.168.1.100 8815

# Start the driver server with 4 local workers (for single-machine distributed execution)
./build/release/distributed_server 0.0.0.0 8815 4

Starting and Registering Standalone Workers

For distributed execution, you can either use local workers (managed by the driver) or register external standalone workers:

Option 1: Start standalone workers in the same process (for testing, debugging and development)

-- Start a local server with driver node
SELECT duckherder_start_local_server(8815);

-- Start standalone workers within the same process
SELECT duckherder_start_standalone_worker(8816);
SELECT duckherder_start_standalone_worker(8817);
SELECT duckherder_start_standalone_worker(8818);

-- Check how many workers are registered
SELECT duckherder_get_worker_count();  -- Returns: 3

Option 2: Register external workers (for multi-machine setups)

-- Start a local server with driver node
SELECT duckherder_start_local_server(8815);

-- Register external worker nodes running on different machines
-- Syntax: duckherder_register_worker(worker_id, location)
SELECT duckherder_register_worker('worker-1', 'grpc://192.168.1.101:8816');
SELECT duckherder_register_worker('worker-2', 'grpc://192.168.1.102:8816');
SELECT duckherder_register_worker('worker-3', 'grpc://192.168.1.103:8816');

-- Verify workers are registered
SELECT duckherder_get_worker_count();  -- Returns: 3

Option 3: Using the standalone distributed_server executable with local workers

# Start the server with 4 automatically managed local workers
./build/release/distributed_server 0.0.0.0 8815 4

Complete Multi-Machine Setup Example

On the driver machine (192.168.1.100):

# Start the standalone driver server
./build/release/distributed_server 0.0.0.0 8815

On each worker machine, run:

# Worker 1 (192.168.1.101)
./build/release/distributed_worker 0.0.0.0 8816 worker-1

# Worker 2 (192.168.1.102)
./build/release/distributed_worker 0.0.0.0 8816 worker-2

# Worker 3 (192.168.1.103)
./build/release/distributed_worker 0.0.0.0 8816 worker-3

From a client machine, register nodes and connect:

-- Register the driver node
SELECT duckherder_register_or_replace_driver('driver-main', 'grpc://192.168.1.100:8815');

-- Register each worker with its network location
SELECT duckherder_register_worker('worker-1', 'grpc://192.168.1.101:8816');
SELECT duckherder_register_worker('worker-2', 'grpc://192.168.1.102:8816');
SELECT duckherder_register_worker('worker-3', 'grpc://192.168.1.103:8816');

-- Confirm registration
SELECT duckherder_get_worker_count();  -- Returns: 3

Working with Remote Tables

-- Create a table.
CREATE TABLE dh.users (
    id INTEGER,
    name VARCHAR,
    email VARCHAR,
    created_at TIMESTAMP
);

-- Insert data.
INSERT INTO dh.users VALUES 
    (1, 'Alice', 'alice@example.com', '2024-01-15 10:30:00'),
    (2, 'Bob', 'bob@example.com', '2024-01-16 14:20:00'),
    (3, 'Charlie', 'charlie@example.com', '2024-01-17 09:15:00');

-- Create an index.
CREATE INDEX idx_users_email ON dh.users(email);

-- Simple SELECT.
SELECT * FROM dh.users WHERE id > 1;

-- Distributed query, which automatically gets parallelized.
SELECT * FROM dh.large_table WHERE value > 1000;

-- Drop a table.
DROP TABLE dh.users;

Query Execution Monitoring and Statistics

duckherder provides built-in stats (current not persisted anywhere) for executed queries, including start timestamp, executime duration and disribution mode.

D SELECT * FROM duckherder_get_query_execution_stats();
┌───────────────────────────────────────┬────────────────┬────────────────┬───────────────────┬──────────────────┬─────────────────────┬────────────────────────────┐
│                  sql                  │ execution_mode │ merge_strategy │ query_duration_ms │ num_workers_used │ num_tasks_generated │    execution_start_time    │
│                varchar                │    varchar     │    varchar     │       int64       │      int64       │        int64        │         timestamp          │
├───────────────────────────────────────┼────────────────┼────────────────┼───────────────────┼──────────────────┼─────────────────────┼────────────────────────────┤
│ SELECT * FROM distributed_basic_table │ DELEGATED      │ CONCATENATE    │                17 │                4 │                   1 │ 2025-11-17 02:10:27.482089 │
│ SELECT * FROM distributed_basic_table │ DELEGATED      │ CONCATENATE    │                 1 │                4 │                   1 │ 2025-11-17 02:10:30.436037 │
│ SELECT * FROM distributed_basic_table │ DELEGATED      │ CONCATENATE    │                 2 │                4 │                   1 │ 2025-11-17 02:10:33.675752 │
│ SELECT * FROM distributed_basic_table │ DELEGATED      │ CONCATENATE    │                 1 │                4 │                   1 │ 2025-11-17 02:10:36.992988 │
└───────────────────────────────────────┴────────────────┴────────────────┴───────────────────┴──────────────────┴─────────────────────┴────────────────────────────┘

Clear Query Statistics

-- Clear all recorded query history and statistics
SELECT duckherder_clear_query_recorder_stats();

-- Verify the history is cleared
SELECT COUNT(*) FROM duckherder_get_query_history();  -- Returns: 0

Current Limitations and Unsupported Features

Duckherder is still experimental. Some features are completely unsupported, while some SQL features are supported only through local execution on the Driver/Control Node.

Not Implemented

  • Remote view discovery after attaching to a catalog that already contains views.
  • CREATE SEQUENCE and remote sequence discovery.
  • ALTER SCHEMA.
  • ALTER TABLE variants other than add/drop/rename column, rename table, change column type, set a default, set/drop NOT NULL, and add a primary key. Adding other constraint types or dropping constraints is not implemented by DuckDB's ALTER path. Foreign keys declared by CREATE TABLE are supported.
  • Full Arrow transport for UNION, BIT, BIGNUM, and non-string dictionary types. Unsupported values may be converted to VARCHAR or rejected.
  • Persistence for the default in-memory object-storage database and server restart recovery.
  • Authentication, authorization, connection pooling, automatic worker replacement, dynamic worker scaling, query result caching, and query resource-consumption tracking.

Supported, but Not Distributed

Worker execution currently accepts only a single table scan with projections, filters, and limited aggregation or GROUP BY. The following features fall back to execution on the Driver/Control Node instead of running across workers:

  • joins and multi-table plans;
  • DISTINCT;
  • ORDER BY;
  • LIMIT and OFFSET;
  • window functions;
  • common table expressions, including recursive CTEs;
  • set operations such as UNION, INTERSECT, and EXCEPT; and
  • all DDL and DML writes.

Experimental or Incomplete

  • Distributed aggregate and GROUP BY finalization uses output-column-name heuristics. Distributed AVG can be incorrect when partitions contain different row counts, and aliases can prevent recognition of COUNT, SUM, MIN, MAX, or AVG.
  • Initial catalog discovery loads schemas, enum types, and tables only. Existing indexes, macros, views, sequences, and other catalog entries are not restored into the client metadata cache after a new ATTACH.
  • Basic BEGIN, COMMIT, and ROLLBACK forwarding is supported, but transaction outcomes are not durable across server restarts. Object-storage snapshot coordination and propagation are still work in progress.
  • DML RETURNING results are materialized and serialized into one protobuf response, so large result sets can use substantial memory.
  • A failed worker task is not retried or reassigned automatically. Replacing a driver also does not synchronize the workers registered with the previous driver.
  • Query execution statistics are held in memory and are not persisted.

See the Roadmap for planned work.

Roadmap

DuckDB compatibility gaps listed below are also tracked by test/configs/duckherder_missing_features.json.

Table and Index Operations

  • Create/drop table
  • Create/drop index
  • Automatic index maintenance on remote DML
  • Basic table schema updates
  • Nested table schema updates
  • Table options and column statistics updates
  • Create/drop view
  • Restore existing views during initial catalog discovery
  • Create/drop schema
  • Create/drop sequence
  • Alter schema and catalog ownership
  • Concurrent catalog DDL/drop semantics
  • DuckDB-compatible eager primary-key and unique-constraint checking

Data Type Support

  • Primitive type support
  • List type support for primitive elements
  • Nested list type support
  • Basic map type support
  • Nested map NULL semantics
  • Struct type support
  • Lossless HUGEINT and UHUGEINT Arrow transport
  • Union type support
  • BIGNUM type support
  • GEOMETRY type and statistics support
  • VARIANT type support

Distributed Query Support

  • Intelligent query partitioning
  • Natural parallelism analysis
  • Row group-aligned execution
  • Range-based partitioning
  • LIMIT pushdown compatibility
  • Aggregation pushdown (infrastructure ready)
  • AVG, BIT_XOR, and nested-list aggregate compatibility
  • Correct GROUP BY distributed finalization
  • JOIN optimization (broadcast/co-partition)
  • ORDER BY support (distributed sort)
  • DISTINCT, LIMIT/OFFSET, window, CTE, and set-operation support
  • TIME and GEOMETRY optimizer statistics
  • Driver collect partition and execution stats

Multi-Client Support

  • Query server authentication and authorization
  • Multiple DuckDB instance support on servers
  • Connection pooling

Full Write Support

  • Persist server-side database file
  • Recover DuckDB instance via database file
  • Basic transaction forwarding
  • INSERT OR REPLACE ... RETURNING NOTHING compatibility
  • Durable transaction and snapshot recovery

Additional Features

  • Query timing statistics
  • Comprehensive execution logging
  • Support official extension install and load
  • Query resource consumption tracking
  • Support community extension install and load
  • Dynamic worker scaling
  • Query result caching
  • COMMENT and pg_description compatibility
  • Hive-partition path escaping for COPY
  • Util function to register driver node and worker nodes
  • Util function to register or replace the driver node

Contributing

See CONTRIBUTING.md for development guidelines.

License

See LICENSE for license information.

About

Distributed execution for duckdb queries.

Resources

Contributing

Stars

142 stars

Watchers

3 watching

Forks

Releases

Packages

Contributors

Languages