Cluster database engine provides real-time access to the tables of a database on a cluster from the server configuration. It is the named-cluster counterpart of the Remote database engine, exactly as the cluster table function relates to the remote table function.
The list of tables and their structure are fetched from the cluster on demand, so the database always reflects its current state. Each table is exposed as a Distributed storage over the named cluster, which forwards SELECT and INSERT queries to it.
This is handy for federating several ClickHouse clusters or for plugging a whole cluster into clickhouse-local or another cluster without spelling out its addresses: the cluster is defined once in the configuration, complete with per-replica credentials, secure connections and the inter-server secret.
Creating a database
cluster_name— The name of a cluster from the server configuration (see Clusters), as in theDistributedtable engine. Macros such as{cluster}are supported and expanded on every access.database— The name of the database on the cluster.
Remote database engine, the Cluster engine takes no credential arguments and stores no secrets: connections use the per-replica settings of the cluster configuration (user, password, secure connections, compression, the inter-server secret). The cluster is re-resolved from the configuration on every access, so the database follows configuration reloads and cluster auto-discovery, like a Distributed table does. The cluster must exist when the database is created; if it later disappears from the configuration, the server still starts and the database reports the missing cluster until the configuration brings it back.
A replica that points to the current server is treated as a local shard: SELECT and INSERT are executed directly under the current user — who therefore needs the corresponding privileges on the underlying database and its tables — and the configured cluster credentials are used only for genuinely remote replicas. If the local replica of a shard does not have the database or a table, the lookup falls back to the remote replicas of the shard, like a Distributed table does.
When the cluster has several shards, each proxy table reads from all of them, but the metadata — the list of the tables and their structure — is taken from an arbitrary shard (a local one is preferred), just like the cluster table function does, so that a listing costs a single query instead of one per shard. The shards of a cluster are therefore expected to serve the same set of tables; a table that only some of them have is served by a proxy whose queries then fail on the shards that do not have it. An INSERT into a table of a multi-shard database sends each row to a random shard (the proxy Distributed tables carry an implicit rand() sharding key, respecting the configured shard weights); to pin the shard for a query, set insert_shard_id. The implicit key only distributes the inserted rows: for reading, the table behaves like a Distributed table without a sharding key (in particular, optimize_skip_unused_shards and force_optimize_skip_unused_shards do not treat it as a shard-pruning key).
In particular, a table that exists only on a non-resolving shard does not appear in SHOW TABLES or EXISTS TABLE, and resolving it by name reports UNKNOWN_TABLE; it has no proxy table to query. A table that exists on the resolving shard but is absent from another shard is visible, but proxy queries can fail on that other shard.
No equivalent re-executable definition exists for a Cluster proxy table, so SHOW CREATE TABLE reports the THERE_IS_NO_QUERY error. This includes the following cases:
- The table is currently served through the remote-replica fallback described above (the local replica of a shard does not have the database or the table). A
Distributedtable over the whole cluster performs no such fallback, and the fallback subset of the cluster has no name of its own. - The database cluster name contains macros.
Clusterexpands those macros on every access, whileDistributedexpands them only when the table is created, so either the macro expression or its current expansion would recreate different behavior after a configuration change. - A configuration reload can change a one-shard cluster to a multi-shard cluster. The live proxy then acquires an implicit
rand()key to distributeINSERTrows, while a table serialized before the reload has no key. Serializing the key unconditionally is not equivalent either, because it would make a standaloneDistributedtable use it for read shard pruning;CREATE TABLEhas no way to express an insert-only sharding key.
Notes
- The engine is a read-through view of the cluster:
CREATE TABLE,DROP TABLE,ALTERand similar DDL statements against theClusterdatabase are not supported. Manage the schema on the cluster directly, e.g. withON CLUSTERDDL queries. - Access rights are enforced on the remote servers for the users configured in the cluster definition (or for the current user when the cluster uses an inter-server secret), and locally by the usual privileges on the database and its tables.
- The behavioral details of the
Remotedatabase engine — the visibility rules for the tables of a local shard, the completeness of the listing, error reporting for an unavailable cluster, and chains of proxy databases — apply to theClusterengine as well; see the notes there.
Example
Create aCluster database that points to the default database of the cluster test_shard_localhost from the server configuration and use it: