spq_plugin_v2:

分支1Tags0
文件最后提交记录最后更新时间
9 个月前
9 个月前
9 个月前
1 天前
9 个月前
9 个月前
9 个月前
9 个月前
9 个月前
1 个月前
9 个月前
9 个月前
9 个月前
18 天前
9 个月前
9 个月前
9 个月前
9 个月前
2 个月前
2 个月前
9 个月前

Citus 数据库是 100% 开源的。

了解 Citus 12.0 版本博客Citus 更新页面 中的新内容。


Citus Banner

最新文档 Stack Overflow Slack 代码覆盖率 Twitter

Citus Deb 软件包 Citus Rpm 软件包

什么是 Citus?

Citus 是一个 PostgreSQL 扩展,它将 Postgres 转变为分布式数据库——让你能够在任何规模下实现高性能。

借助 Citus,你可以为 PostgreSQL 数据库扩展新的强大功能:

  • 分布式表 在 PostgreSQL 节点集群中进行分片,以整合它们的 CPU、内存、存储和 I/O 容量。
  • 引用表 复制到所有节点,用于分布式表的连接和外键操作,并实现最高的读取性能。
  • 分布式查询引擎 在集群中路由和并行处理分布式表上的 SELECT、DML 和其他操作。
  • 列式存储 压缩数据,加速扫描,并支持对常规表和分布式表的快速投影。
  • 从任何节点查询 使你能够利用集群的全部容量进行分布式查询

你可以使用这些 Citus 强大功能,让你的 Postgres 数据库在单个 Citus 节点上具备横向扩展能力。或者,你可以构建一个大型集群,能够处理高事务吞吐量(尤其是在多租户应用中),运行快速分析查询,并处理大量时间序列物联网数据以进行实时分析。当你的数据大小和容量增长时,你可以轻松地向集群添加更多工作节点并重新平衡分片。

我们在 SIGMOD '21 会议上发表的论文 Citus: Distributed PostgreSQL for Data-Intensive Applications 更详细地介绍了 Citus 是什么、它如何工作以及为什么这样工作。

Citus 从单个节点横向扩展

由于 Citus 是 Postgres 的扩展,你可以将 Citus 与最新的 Postgres 版本一起使用。并且 Citus 能与你已经熟悉的 PostgreSQL 工具和扩展无缝协作。

为何选择 Citus?

开发人员选择 Citus 主要基于以下两个原因:

  1. 您的应用程序已超出单个 PostgreSQL 节点的承载能力

    随着数据规模和数据量的不断增长,单个 PostgreSQL 节点可能会开始出现各种性能和可扩展性问题。例如:CPU 利用率过高和 I/O 等待时间过长导致查询变慢,SQL 查询返回内存不足错误,自动清理(autovacuum)无法跟上并导致表膨胀等。

    借助 Citus,您可以对表进行分布式处理并可选择进行压缩,从而确保始终有足够的内存、CPU 和 I/O 容量,以实现在大规模数据下的高性能。分布式查询引擎能够高效地在集群中路由事务,同时在所有核心上并行处理分析查询和批处理操作。此外,您仍然可以使用您熟悉和喜爱的 PostgreSQL 功能与工具。

  2. PostgreSQL 具备其他系统所不具备的能力

    市面上有许多旨在横向扩展的数据处理系统,但很少有系统能像 PostgreSQL 一样拥有如此多强大的功能,包括:高级连接和子查询、用户定义函数、更新/删除/插入更新(upsert)、约束和外键、强大的扩展(如 PostGIS、HyperLogLog)、多种类型的索引、时间分区以及完善的 JSON 支持。

    Citus 使 PostgreSQL 最强大的功能能够在任何规模下发挥作用,让您能够在单一数据库系统上处理复杂的数据密集型工作负载。

快速入门

开始使用 Citus 最快捷的方式是使用云中的 Azure Cosmos DB for PostgreSQL 托管服务,或者在本地设置 Citus

Azure 上的 Citus 托管服务

您可以通过 Azure Cosmos DB for PostgreSQL 门户 在几分钟内获得一个完全托管的 Citus 集群。Azure 将为您所有的服务器管理备份、通过自动故障转移实现的高可用性、软件更新、监控等。要开始在 Azure 上使用 Citus,请参阅 Azure Cosmos DB for PostgreSQL 快速入门

使用 Docker 运行 Citus

最小的 Citus 集群是一个带有 Citus 扩展的单个 PostgreSQL 节点,这意味着您可以通过运行单个 Docker 容器来试用 Citus。

# run PostgreSQL with Citus on port 5500
docker run -d --name citus -p 5500:5432 -e POSTGRES_PASSWORD=mypassword citusdata/citus

# connect using psql within the Docker container
docker exec -it citus psql -U postgres

# or, connect using local psql
psql -U postgres -d postgres -h localhost -p 5500

在本地安装 Citus

如果您已在本地安装 PostgreSQL,安装 Citus 最简单的方法是使用我们的打包仓库。

在 Ubuntu / Debian 上安装软件包:

curl https://install.citusdata.com/community/deb.sh > add-citus-repo.sh
sudo bash add-citus-repo.sh
sudo apt-get -y install postgresql-15-citus-12.0

在 CentOS / Red Hat 上安装软件包:

curl https://install.citusdata.com/community/rpm.sh > add-citus-repo.sh
sudo bash add-citus-repo.sh
sudo yum install -y citus120_15

要将 Citus 添加到本地 PostgreSQL 数据库,请在 postgresql.conf 中添加以下内容:

shared_preload_libraries = 'citus'

重启 PostgreSQL 后,使用 psql 连接并运行:

CREATE EXTENSION citus;

您现在已准备就绪,可以开始在单节点上使用 Citus 表了。

在多节点上安装 Citus

如果您想设置多节点集群,也可以设置其他带有 Citus 扩展的 PostgreSQL 节点,并将它们添加进来以形成 Citus 集群:

-- before adding the first worker node, tell future worker nodes how to reach the coordinator
SELECT citus_set_coordinator_host('10.0.0.1', 5432);

-- add worker nodes
SELECT citus_add_node('10.0.0.2', 5432);
SELECT citus_add_node('10.0.0.3', 5432);

-- rebalance the shards over the new worker nodes
SELECT rebalance_table_shards();

有关更多详细信息,请参阅我们的多节点 Citus 集群设置文档,其中介绍了在各种操作系统上的安装方法。

使用 Citus

搭建好 Citus 集群后,您就可以开始创建分布式表、引用表并使用列式存储了。

创建分布式表

create_distributed_table UDF 会自动在本地或跨工作节点对表进行分片:

CREATE TABLE events (
  device_id bigint,
  event_id bigserial,
  event_time timestamptz default now(),
  data jsonb not null,
  PRIMARY KEY (device_id, event_id)
);

-- distribute the events table across shards placed locally or on the worker nodes
SELECT create_distributed_table('events', 'device_id');

完成此操作后,针对特定设备 ID 的查询将被高效路由至单个工作节点,而跨设备 ID 的查询则会在集群中并行处理。

-- insert some events
INSERT INTO events (device_id, data)
SELECT s % 100, ('{"measurement":'||random()||'}')::jsonb FROM generate_series(1,1000000) s;

-- get the last 3 events for device 1, routed to a single node
SELECT * FROM events WHERE device_id = 1 ORDER BY event_time DESC, event_id DESC LIMIT 3;
┌───────────┬──────────┬───────────────────────────────┬───────────────────────────────────────┐
│ device_id │ event_id │          event_time           │                 data                  │
├───────────┼──────────┼───────────────────────────────┼───────────────────────────────────────┤
│         119999012021-03-04 16:00:31.189963+00 │ {"measurement": 0.88722643925054}     │
│         119998012021-03-04 16:00:31.189963+00 │ {"measurement": 0.6512231304621992}   │
│         119997012021-03-04 16:00:31.189963+00 │ {"measurement": 0.019368766051897524} │
└───────────┴──────────┴───────────────────────────────┴───────────────────────────────────────┘
(3 rows)

Time: 4.588 ms

-- explain plan for a query that is parallelized across shards, which shows the plan for
-- a query one of the shards and how the aggregation across shards is done
EXPLAIN (VERBOSE ON) SELECT count(*) FROM events;
┌────────────────────────────────────────────────────────────────────────────────────┐
│                                     QUERY PLAN                                     │
├────────────────────────────────────────────────────────────────────────────────────┤
│ Aggregate                                                                          │
│   Output: COALESCE((pg_catalog.sum(remote_scan.count))::bigint, '0'::bigint)       │
│   ->  Custom Scan (Citus Adaptive)                                                 │
│         ...                                                                        │
│         ->  Task                                                                   │
│               Query: SELECT count(*) AS count FROM events_102008 events WHERE true │
│               Node: host=localhost port=5432 dbname=postgres                       │
│               ->  Aggregate                                                        │
│                     ->  Seq Scan on public.events_102008 events                    │
└────────────────────────────────────────────────────────────────────────────────────┘

创建具有共置功能的分布式表

具有相同分布列的分布式表可以进行共置,以实现分布式表之间的高性能分布式连接和外键。 默认情况下,分布式表将根据分布列的类型进行共置,但您可以在 create_distributed_table 中使用 colocate_with 参数显式定义共置。

CREATE TABLE devices (
  device_id bigint primary key,
  device_name text,
  device_type_id int
);
CREATE INDEX ON devices (device_type_id);

-- co-locate the devices table with the events table
SELECT create_distributed_table('devices', 'device_id', colocate_with := 'events');

-- insert device metadata
INSERT INTO devices (device_id, device_name, device_type_id)
SELECT s, 'device-'||s, 55 FROM generate_series(0, 99) s;

-- optionally: make sure the application can only insert events for a known device
ALTER TABLE events ADD CONSTRAINT device_id_fk
FOREIGN KEY (device_id) REFERENCES devices (device_id);

-- get the average measurement across all devices of type 55, parallelized across shards
SELECT avg((data->>'measurement')::double precision)
FROM events JOIN devices USING (device_id)
WHERE device_type_id = 55;

┌────────────────────┐
│        avg         │
├────────────────────┤
│ 0.5000191877513974 │
└────────────────────┘
(1 row)

Time: 209.961 ms

共置还能帮助你扩展 INSERT..SELECT存储过程分布式事务

不间断应用的情况下分布表

你们中的一些人可能已经开始使用 Postgres,并决定在应用程序使用表的过程中稍后再分布表。在这种情况下,你希望避免读写操作的停机时间。create_distributed_table 命令会阻塞表上的写入操作(例如 DML 命令),直到命令完成。相反,使用 create_distributed_table_concurrently 命令,即使在命令执行期间,你的应用程序也可以继续读写数据。

CREATE TABLE device_logs (
  device_id bigint primary key,
  log text
);

-- insert device logs
INSERT INTO device_logs (device_id, log)
SELECT s, 'device log:'||s FROM generate_series(0, 99) s;

-- convert device_logs into a distributed table without interrupting the application
SELECT create_distributed_table_concurrently('device_logs', 'device_id', colocate_with := 'devices');


-- get the count of the logs, parallelized across shards
SELECT count(*) FROM device_logs;

┌───────┐
│ count │
├───────┤
│   100 │
└───────┘
(1 row)

Time: 48.734 ms

创建引用表

当您需要不包含分布列的快速连接或外键时,可以使用 create_reference_table 在集群的所有节点上复制表。

CREATE TABLE device_types (
  device_type_id int primary key,
  device_type_name text not null unique
);

-- replicate the table across all nodes to enable foreign keys and joins on any column
SELECT create_reference_table('device_types');

-- insert a device type
INSERT INTO device_types (device_type_id, device_type_name) VALUES (55, 'laptop');

-- optionally: make sure the application can only insert devices with known types
ALTER TABLE devices ADD CONSTRAINT device_type_fk
FOREIGN KEY (device_type_id) REFERENCES device_types (device_type_id);

-- get the last 3 events for devices whose type name starts with laptop, parallelized across shards
SELECT device_id, event_time, data->>'measurement' AS value, device_name, device_type_name
FROM events JOIN devices USING (device_id) JOIN device_types USING (device_type_id)
WHERE device_type_name LIKE 'laptop%' ORDER BY event_time DESC LIMIT 3;

┌───────────┬───────────────────────────────┬─────────────────────┬─────────────┬──────────────────┐
│ device_id │          event_time           │        value        │ device_name │ device_type_name │
├───────────┼───────────────────────────────┼─────────────────────┼─────────────┼──────────────────┤
│        602021-03-04 16:00:31.189963+000.28902084163415864 │ device-60   │ laptop           │
│         82021-03-04 16:00:31.189963+000.8723803076285073  │ device-8    │ laptop           │
│        202021-03-04 16:00:31.189963+000.8177634801548557  │ device-20   │ laptop           │
└───────────┴───────────────────────────────┴─────────────────────┴─────────────┴──────────────────┘
(3 rows)

Time: 146.063 ms

参考表能让您扩展复杂的数据模型,并充分利用关系型数据库的特性。

使用列存储创建表

要在您的 PostgreSQL 数据库中使用列存储,只需在 CREATE TABLE 语句中添加 USING columnar,您的数据就会通过列存储访问方式自动进行压缩。

CREATE TABLE events_columnar (
  device_id bigint,
  event_id bigserial,
  event_time timestamptz default now(),
  data jsonb not null
)
USING columnar;

-- insert some data
INSERT INTO events_columnar (device_id, data)
SELECT d, '{"hello":"columnar"}' FROM generate_series(1,10000000) d;

-- create a row-based table to compare
CREATE TABLE events_row AS SELECT * FROM events_columnar;

-- see the huge size difference!
\d+
                                          List of relations
┌────────┬──────────────────────────────┬──────────┬───────┬─────────────┬────────────┬─────────────┐
│ Schema │             Name             │   Type   │ Owner │ Persistence │    Size    │ Description │
├────────┼──────────────────────────────┼──────────┼───────┼─────────────┼────────────┼─────────────┤
│ public │ events_columnar              │ table    │ marco │ permanent   │ 25 MB      │             │
│ public │ events_row                   │ table    │ marco │ permanent   │ 651 MB     │             │
└────────┴──────────────────────────────┴──────────┴───────┴─────────────┴────────────┴─────────────┘
(2 rows)

你可以单独使用列式存储,也可以在分布式表中结合使用,以同时享受压缩和分布式查询引擎带来的优势。

使用列式存储时,应仅通过 COPYINSERT..SELECT 批量加载数据,以实现良好的压缩效果。目前,列式表暂不支持更新、删除操作以及外键。不过,你可以使用分区表,让较新的分区采用行式存储,而较旧的分区则使用列式存储进行压缩。

要了解有关列式存储的更多信息,请查看 列式存储 README

基于模式的分片

自 Citus 12.0 起可用,基于模式的分片 采用共享数据库、独立模式的模型,模式成为数据库内的逻辑分片。多租户应用可以为每个租户使用一个模式,以便轻松地沿着租户维度进行分片。无需更改查询,应用通常只需稍作修改,在切换租户时设置正确的 search_path 即可。基于模式的分片是微服务以及独立软件开发商(ISV)部署无法进行行式分片所需更改的应用程序的理想解决方案。

创建分布式模式

你可以通过调用 citus_schema_distribute 将现有模式转换为分布式模式:

SELECT citus_schema_distribute('user_service');

或者,你可以设置 citus.enable_schema_based_sharding,使所有新创建的模式自动转换为分布式模式:

SET citus.enable_schema_based_sharding TO ON;

CREATE SCHEMA AUTHORIZATION user_service;
CREATE SCHEMA AUTHORIZATION time_service;
CREATE SCHEMA AUTHORIZATION ping_service;

运行查询

查询将根据 search_path 正确路由到相应的模式,或者通过在查询中显式使用模式名称来路由。

对于微服务,您需要为每个服务创建一个与模式名称匹配的 USER,因此默认的 search_path 将包含该模式名称。连接后,用户查询将自动路由,无需对微服务进行任何更改。

CREATE USER user_service;
CREATE SCHEMA AUTHORIZATION user_service;

对于典型的多租户应用程序,你需要在应用中将搜索路径设置为租户的 schema 名称:

SET search_path = tenant_name, public;

高可用设置

作为 PostgreSQL 最受欢迎的高可用解决方案之一,Patroni 3.0 对 Citus 10.0 及以上版本提供了一流支持。此外,自 Citus 11.2 版本起,其针对 Patroni 中的节点切换进行了优化,使切换过程更加顺畅。

Citus 集群的 patronictl list 输出示例:

postgres@coord1:~$ patronictl list demo
+ Citus cluster: demo ----------+--------------+---------+----+-----------+
| Group | Member  | Host        | Role         | State   | TL | Lag in MB |
+-------+---------+-------------+--------------+---------+----+-----------+
|     0 | coord1  | 172.27.0.10 | Replica      | running |  1 |         0 |
|     0 | coord2  | 172.27.0.6  | Sync Standby | running |  1 |         0 |
|     0 | coord3  | 172.27.0.4  | Leader       | running |  1 |           |
|     1 | work1-1 | 172.27.0.8  | Sync Standby | running |  1 |         0 |
|     1 | work1-2 | 172.27.0.2  | Leader       | running |  1 |           |
|     2 | work2-1 | 172.27.0.5  | Sync Standby | running |  1 |         0 |
|     2 | work2-2 | 172.27.0.7  | Leader       | running |  1 |           |
+-------+---------+-------------+--------------+---------+----+-----------+

文档

如果您已准备好开始使用 Citus 或想了解更多信息,建议阅读 Citus 开源文档。如果您在 Azure 上使用 Citus,那么 Azure Cosmos DB for PostgreSQL 是您的起点。

我们的 Citus 文档包含全面的用例指南,介绍如何构建 多租户 SaaS 应用程序实时分析仪表板,或处理 时间序列数据

架构

Citus 数据库集群从单个 PostgreSQL 节点开始,通过添加工作节点扩展为集群。在 Citus 集群中,应用程序连接的原始节点称为协调节点。Citus 协调节点包含分布式表和引用表的元数据,以及常规(本地)表、序列和其他数据库对象(例如外部表)。

分布式表中的数据存储在“分片”中,这些分片实际上是工作节点上的常规 PostgreSQL 表。在协调节点上查询分布式表时,Citus 会向工作节点发送常规 SQL 查询。这样,所有常见的 PostgreSQL 优化和扩展都可以自动与 Citus 一起使用。

Citus 架构

当您发送的查询中,所有(共置的)分布式表在分布列上具有相同的筛选条件时,Citus 会自动检测到这一点,并将整个查询发送到存储数据的工作节点。这样,支持任意复杂的查询,且路由开销最小,这对于扩展事务工作负载特别有用。如果查询没有特定的筛选条件,则会并行查询每个分片,这在分析工作负载中尤其有用。Citus 分布式执行器具有自适应性,旨在在高并发情况下在同一系统上同时处理这两种查询类型,从而支持大规模混合工作负载。

分布式表和引用表的架构及元数据会自动同步到集群中的所有节点。这样,您可以连接到任何节点来运行分布式查询。架构更改和集群管理仍需通过协调节点进行。

何时使用 Citus

Citus 具备独特的能力,可同时扩展分析型和事务型工作负载,支持高达 PB 级的数据量。Citus 常见的使用场景如下:

  • 面向客户的分析仪表板: Citus 让您能够构建分析仪表板,该仪表板可同时在数据库中摄入和处理大量数据,即使面对大量并发用户,也能提供亚秒级响应时间。

    Citus 先进的并行分布式查询引擎,结合 PostgreSQL 的特性(如 数组类型JSONB横向连接)以及 HyperLogLogTopN 等扩展,使您能够构建响应迅速的分析仪表板,无论您有多少客户或数据量多大。

    实时分析用户案例:AlgoliaHeap

  • 时间序列数据: Citus 使您能够处理和分析海量时间序列数据。最大的 Citus 集群可存储超过 1 PB 的时间序列数据,每天摄入数 TB 的数据。

    Citus 与 Postgres 表分区 无缝集成,并具有 按时间分区的内置函数,这可以加快时间序列表上的查询和写入速度。您可以利用 Citus 的并行分布式查询引擎进行快速分析查询,并使用内置的 列式存储 来压缩旧分区。

    用户案例:MixRankWindows 团队

  • 软件即服务 (SaaS) 应用: SaaS 和其他多租户应用需要能够随着租户/客户数量的增长而扩展其数据库。Citus 使您能够按租户维度透明地分片复杂的数据模型,从而使您的数据库能够随着业务增长而扩展。

    通过沿租户 ID 列分布表并将同一租户的数据共置,Citus 可以水平扩展复杂的(租户范围内的)查询、事务和外键图。引用表和分布式 DDL 命令使数据库管理比手动分片更加轻松。此外,您还拥有内置的分布式查询引擎,可在数据库内进行跨租户分析。

    多租户 SaaS 用户案例:CopperSalesloftConvertFlow

  • 微服务:Citus 支持基于模式的分片,允许将常规数据库模式分布在多台机器上。这种分片方法非常适合典型的微服务架构,在这种架构中,存储完全由服务拥有,因此不能与其他租户共享相同的模式定义。Citus 允许跨服务分布水平可扩展的状态,解决了微服务的 主要问题之一

  • 地理空间: 由于强大的 PostGIS 扩展为 Postgres 添加了对地理对象的支持,许多人在 Postgres 之上运行空间/GIS 应用。而且,由于空间位置信息已成为我们日常生活的一部分,地理空间应用比以往任何时候都多。当您的 Postgres 数据库需要扩展以处理增加的工作负载时,Citus 是一个很好的选择。

    地理空间用户案例:赫尔辛基地区交通管理局 (HSL)MobilityDB

需要帮助?

贡献

Citus 建立在开源之上,也属于开源,我们欢迎您的贡献。CONTRIBUTING.md 文件说明了如何开始开发 Citus 扩展本身以及我们的代码质量准则。

行为准则

本项目采用了 Microsoft 开源行为准则。 欲了解更多信息,请参阅 行为准则常见问题,或通过 opencode@microsoft.com 联系我们以获取任何其他问题或意见。

保持联系


Copyright © Citus Data, Inc.