TimescaleDB 集成
TimescaleDB 集成
TimescaleDB 是一个开源数据库,旨在使 SQL 可扩展地处理时序数据。它基于 PostgreSQL 数据库,并以扩展的形式实现。
Debezium PostgreSQL 连接器可以从 TimescaleDB 中捕获数据变更。标准的 PostgreSQL 连接器从数据库中读取原始数据。然后,您可以使用 io.debezium.connector.postgresql.transforms.timescaledb.TimescaleDb 转换来处理原始数据、执行逻辑路由并添加相关元数据。
安装
- 按照 TimescaleDB 文档中的说明安装 TimescaleDB。
- 按照 Debezium 安装指南中的说明安装 Debezium PostgreSQL 连接器。
- 配置 TimescaleDB 并部署连接器。
工作原理
Debezium 可以从以下 TimescaleDB 功能中捕获事件:
- 超表(Hypertables)
- 连续聚合(Continuous aggregates)
- 压缩(Compression)
这三个功能在内部是相互依赖的。每个功能都建立在 PostgreSQL 将数据存储到表中的功能之上。Debezium 对这三个功能的支持程度各不相同。
SMT(单消息转换)需要访问 TimescaleDB 元数据。由于 SMT 无法在连接器级别访问数据库配置,因此您必须为该转换显式定义配置元数据。
超表
超表是一种用于存储时序数据的逻辑表。数据会根据定义的时间边界列进行分块(分区)。TimescaleDB 会在其内部模式中创建一个或多个物理表,每个表代表一个分块。默认情况下,连接器会从每个分块表中捕获变更,并将这些变更流式传输到与各分块对应的各个主题中。TimescaleDB 转换会从各个主题中重新组装数据,然后将重新组装后的数据流式传输到单个主题中。
该转换可以访问 TimescaleDB 元数据以获取分块与超表的映射关系。
该转换会将捕获的事件从其特定分块主题重新路由到单个逻辑主题,主题按照以下模式命名:
<prefix>.<hypertable-schema-name>.<hypertable-name>该转换会向事件添加以下标头:
__debezium_timescaledb_chunk_table存储事件数据的物理表的名称。
__debezium_timescaledb_chunk_schema该物理表所属的模式名称。
示例:从超表流式传输数据
以下示例展示了在 public 模式中创建 conditions 超表的 SQL 命令:
CREATE TABLE conditions (time TIMESTAMPTZ NOT NULL, location TEXT NOT NULL, temperature DOUBLE PRECISION NULL, humidity DOUBLE PRECISION NULL);
SELECT create_hypertable('conditions', 'time');TimescaleDB SMT 会将从超表中捕获的变更事件路由到名为 timescaledb.public.conditions 的主题。该转换会根据你在配置中定义的内容,为事件消息添加相应的头部信息。例如:
__debezium_timescaledb_chunk_table: _hyper_1_1_chunk
__debezium_timescaledb_chunk_schema: _timescaledb_internal连续聚合
连续聚合会对存储在超表中的数据自动进行统计计算。聚合视图由其自身的超表支撑,而该超表又由一组 PostgreSQL 表支撑。这些聚合可以自动或手动重新计算。聚合重新计算后,新的值会存储在超表中,从而可以被捕获并流式传输。聚合中的数据会根据其存储所在的分块(chunk)被流式传输到不同的主题。TimescaleDB 转换器会将被流式传输到不同主题的数据重新组装,并将其路由到单个主题。
该转换器可以访问 TimescaleDB 的元数据,以获取分块与超表之间、以及超表与聚合之间的映射关系。
该转换器会将捕获的事件从其各自的分块专属主题重新路由到单个逻辑主题,该主题按照以下模式命名:
<prefix>.<aggregate-schema-name>.<aggregate-name>。该转换器会为事件添加以下头部:
__debezium_timescaledb_hypertable_table存储该连续聚合的超表的名称。
__debezium_timescaledb_hypertable_schema该超表所属架构(schema)的名称。
__debezium_timescaledb_chunk_table存储该连续聚合的物理表的名称。
__debezium_timescaledb_chunk_schema该物理表所属架构的名称。
示例:从连续聚合中流式传输数据
以下示例展示了在 public 架构中创建名为 conditions_summary 的连续聚合的 SQL 命令。
CREATE MATERIALIZED VIEW conditions_summary WITH (timescaledb.continuous) AS
SELECT
location,
time_bucket(INTERVAL '1 hour', time) AS bucket,
AVG(temperature),
MAX(temperature),
MIN(temperature)
FROM conditions
GROUP BY location, bucket;TimescaleDB SMT 会将聚合中捕获的变更事件路由到名为 timescaledb.public.conditions_summary 的主题。该转换会根据配置中定义的内容,为事件消息添加相应的报头。例如:
_debezium_timescaledb_chunk_table: _hyper_2_2_chunk
__debezium_timescaledb_chunk_schema: _timescaledb_internal
__debezium_timescaledb_hypertable_table: _materialized_hypertable_2
__debezium_timescaledb_hypertable_schema: _timescaledb_internal压缩
TimescaleDB SMT 不会对压缩函数做任何特殊处理。压缩的数据块会原样转发到管道中的下一个下游任务,以便根据需要进一步处理。通常情况下,包含压缩数据块的消息会被丢弃,不会被管道中的后续任务处理。
TimescaleDB 配置
Debezium 使用复制槽(replication slot)从 TimescaleDB 和 PostgreSQL 中捕获变更。复制槽以多种消息格式存储数据。通常,最好将 Debezium 配置为使用 pgoutput 解码器从复制槽中读取数据,pgoutput 是 TimescaleDB 实例的默认解码器。
要配置复制槽以支持逻辑解码,请在 postgresql.conf 文件中指定以下设置:
# REPLICATION
wal_level = logical必须重启数据库服务器,设置才能生效。
要配置需要复制的表,必须创建一个发布(publication),如下例所示:
CREATE PUBLICATION dbz_publication FOR ALL TABLES WITH (publish = 'insert, update')你可以像前面的示例那样创建全局发布,也可以为每张表分别创建发布。由于 TimescaleDB 会按需自动创建表,因此强烈建议使用全局发布。
连接器配置
TimescaleDB SMT 的配置方式与 PostgreSQL 连接器相同。为使连接器能够正确处理来自 TimescaleDB 的事件,请在连接器配置中添加以下选项:
"transforms": "timescaledb",
"transforms.timescaledb.type": "io.debezium.connector.postgresql.transforms.timescaledb.TimescaleDb",
"transforms.timescaledb.database.hostname": "timescaledb",
"transforms.timescaledb.database.port": "...",
"transforms.timescaledb.database.user": "...",
"transforms.timescaledb.database.password": "...",
"transforms.timescaledb.database.dbname": "..."连接器配置示例
下面的示例展示了如何配置 PostgreSQL 连接器,以通过逻辑名称 dbserver1 连接到位于 192.168.99.100、端口 5432 上的 TimescaleDB 服务器。通常,你会在一个 JSON 文件中通过设置连接器提供的配置属性来配置 Debezium PostgreSQL 连接器。
你可以选择只为数据库中的一部分模式和表生成事件。此外,你还可以忽略、掩码或截断包含敏感数据的列、超出指定大小的列,或你不需要的列。
{
"name": "timescaledb-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "192.168.99.100",
"database.port": "5432",
"database.user": "postgres",
"database.password": "postgres",
"database.dbname" : "postgres",
"topic.prefix": "dbserver1",
"plugin.name": "pgoutput",
"schema.include.list": "_timescaledb_internal",
"transforms": "timescaledb",
"transforms.timescaledb.type": "io.debezium.connector.postgresql.transforms.timescaledb.TimescaleDb",
"transforms.timescaledb.database.hostname": "timescaledb",
"transforms.timescaledb.database.port": "5432",
"transforms.timescaledb.database.user": "postgres",
"transforms.timescaledb.database.password": "postgres",
"transforms.timescaledb.database.dbname": "postgres"
}
}以下列表说明了上例中各行的用途:
"name": "timescaledb-connector",
向 Kafka Connect 服务注册时,该连接器的名称。
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
此 PostgreSQL 连接器类的名称。
"database.hostname": "192.168.99.100",
TimescaleDB 服务器的地址。
"database.port": "5432",
TimescaleDB 服务器的端口号。
"database.user": "postgres",
TimescaleDB 用户的名称。
"database.password": "postgres",
TimescaleDB 的密码。
"database.dbname" : "postgres",
要连接的 TimescaleDB 数据库的名称。
"topic.prefix": "dbserver1",
TimescaleDB 服务器或集群的主题前缀。该前缀构成一个命名空间,用于连接器写入的所有 Kafka 主题的名称、Kafka Connect 架构名称,以及在使用 Avro 转换器时相应 Avro 架构的命名空间。
"plugin.name": "pgoutput",
表示使用 pgoutput 逻辑解码插件。
"schema.include.list": "_timescaledb_internal",
包含 TimescaleDB 物理表的所有架构列表。
"transforms": "timescaledb",
启用 SMT 以处理原始 TimescaleDB 事件。
"transforms.timescaledb.type": "io.debezium.connector.postgresql.transforms.timescaledb.TimescaleDb",
启用 SMT 以处理原始 TimescaleDB 事件。
"transforms.timescaledb.database.hostname": "timescaledb",
此行及后续各行向 SMT 提供 TimescaleDB 连接信息,包括主机名、端口号、数据库用户、数据库密码和数据库名称。这些值必须与前文 database.hostname、database.port、database.user、database.password 和 database.dbname 字段中的值一致。
配置选项
下表列出了为 TimescaleDB 集成 SMT 可设置的配置选项。
| 属性 | 默认值 | 描述 | |
|---|---|---|---|
database.hostname | 无默认值 | TimescaleDB 数据库服务器的 IP 地址或主机名。 | |
database.port | 5432 | TimescaleDB 数据库服务器的整数端口号。 | |
database.user | 无默认值 | 连接 TimescaleDB 数据库服务器时所使用的 TimescaleDB 数据库用户名。 | |
database.password | 无默认值 | 连接 TimescaleDB 数据库服务器时要使用的密码。 | |
database.dbname | 无默认值 | 要从中流式传输变更的 TimescaleDB 数据库名称。 | |
schema.list | _timescaledb_internal | 以逗号分隔的架构(schema)名称列表,这些架构中包含 TimescaleDB 原始(内部)数据表。SMT 仅处理源自该列表中某个架构的变更。 | |
target.topic.prefix | timescaledb | TimescaleDB 事件被路由到的主题的命名空间(前缀)。SMT 会将消息路由到名为 `.._<hypertable | aggregate>_` 的主题中。 |
表 1. TimescaleDB 集成 SMT(TimescaleDB)配置选项
评论
登录后参与评论
KnowForge