About this project

Citus is a PostgreSQL extension that transforms a single Postgres database into a distributed one, letting you scale out across multiple nodes while continuing to use standard PostgreSQL tools and extensions. Core capabilities described in the README: - Distributed tables: tables are sharded across a cluster of PostgreSQL nodes, combining CPU, memory, storage and I/O capacity. The create_distributed_table function shards a table locally or across worker nodes. - Reference tables: replicated to all nodes to support joins and foreign keys that do not include the distribution column, and to improve read performance. - Distributed query engine: routes and parallelizes SELECT, DML and other operations across the cluster. Queries filtered on the distribution column are routed to a single worker; other queries are parallelized across shards. - Columnar storage: a USING columnar table access method that compresses data, speeds up scans and supports fast projections, usable on regular or distributed tables. Updates, deletes and foreign keys are not supported on columnar tables; batch loading via COPY or INSERT..SELECT is recommended. - Query from any node: distributed queries can be issued from any node in the cluster, while schema changes and cluster administration still go through the coordinator. - Schema-based sharding (since Citus 12.0): a shared-database, separate-schema model where each schema acts as a logical shard, useful for multi-tenant apps and microservices without query changes. - create_distributed_table_concurrently: converts an existing table to a distributed one without blocking reads and writes. - Co-location: distributed tables sharing a distribution column can be co-located to enable efficient distributed joins, foreign keys, INSERT..SELECT, stored procedures and distributed transactions. Deployment options include a managed service on Azure (Elastic Clusters in Azure Database for PostgreSQL Flexible Server), a Docker image (citusdata/citus), and local packages for Ubuntu/Debian and Red Hat. Multi-node clusters are built by adding worker nodes and rebalancing shards. High availability is supported through Patroni 3.0, which has first-class Citus support. The README positions Citus for applications outgrowing a single PostgreSQL node and for workloads needing PostgreSQL features at scale, such as multi-tenant SaaS, real-time analytics dashboards, time-series and IoT data. Documentation, use-case guides and a SIGMOD '21 paper are referenced for further detail.