选择 打开 改范围 完整检索页
受支持版本: 当前版本 (18) / 17 / 16 / 15 / 14
开发版本: 19 / devel
不受支持的版本: 13 / 12 / 11 / 10
当前 PostgreSQL 版本不在支持生命周期内。
您可以参阅当前版本的对应页面,或其他在上面列出的活跃大版本。

49.6. 逻辑解码输出插件 #

PostgreSQL 源码树中的 contrib/test_decoding 子目录里有一个输出插件示例。

49.6.1. 初始化函数 #

加载输出插件的方式是动态加载共享库,并以输出插件的名称作为库的基本名称。定位该库时使用常规库搜索路径。为了提供所需的输出插件回调,并表明这个库确实是输出插件,库必须提供一个名为_PG_output_plugin_init的函数。此函数会接收一个结构体,需要在其中填写各项操作的回调函数指针。

typedef struct OutputPluginCallbacks
{
    LogicalDecodeStartupCB startup_cb;
    LogicalDecodeBeginCB begin_cb;
    LogicalDecodeChangeCB change_cb;
    LogicalDecodeTruncateCB truncate_cb;
    LogicalDecodeCommitCB commit_cb;
    LogicalDecodeMessageCB message_cb;
    LogicalDecodeFilterByOriginCB filter_by_origin_cb;
    LogicalDecodeShutdownCB shutdown_cb;
} OutputPluginCallbacks;

typedef void (*LogicalOutputPluginInit) (struct OutputPluginCallbacks *cb);

其中,begin_cbchange_cbcommit_cb回调是必需的,而startup_cbfilter_by_origin_cbtruncate_cb,以及shutdown_cb是可选的。如果没有设置truncate_cb,但要解码一个TRUNCATE操作,那么该操作会被忽略。

49.6.2. 能力 #

为了对变更进行解码、格式化和输出,输出插件可以使用后端的大部分常规基础设施,包括调用输出函数。可以对关系进行只读访问,但仅限于以下两种关系:由initdb创建在pg_catalog模式中的关系,或者使用以下方式标记为用户提供的系统目录表的关系:

ALTER TABLE user_catalog_table SET (user_catalog_table = true);
CREATE TABLE another_catalog_table(data text) WITH (user_catalog_table = true);

禁止执行任何会导致分配事务 ID 的操作。这包括向表写入数据、执行 DDL 更改以及调用txid_current()

49.6.3. 输出模式 #

输出插件回调几乎可以用任意格式向消费者传递数据。对于某些用例,例如通过 SQL 查看更改,把数据返回为能够容纳任意数据的数据类型(例如 bytea)会很笨拙。如果输出插件只输出服务器编码中的文本数 据,它可以把 OutputPluginOptions.output_type 设置为 OUTPUT_PLUGIN_TEXTUAL_OUTPUT 而不是 OUTPUT_PLUGIN_BINARY_OUTPUT,并在 启动回调 中声明这一点。在这种情况下,所有数据都必须采用服务器编码,这样才能装入 text datum。断言开启的构建会检查这一点。

49.6.4. 输出插件回调 #

输出插件通过它所提供的各类回调获知正在发生的更改。

并发事务按提交顺序解码,并且只有属于某个特定事务的更改,才会在 begincommit 回调之间被解码。 显式或隐式回滚的事务永远不会被解码。成功的保存点会按照它们在该事务中执行的顺序,被折叠进包含它们的事务中。

Note

只有已经安全刷写到磁盘的事务才会被解码。这可能导致 COMMIT 在紧随其后的 pg_logical_slot_get_changes() 调用中不会立即被解 码,当 synchronous_commit 被设置为 off 时尤其如此。

49.6.4.1. 启动回调 #

可选的startup_cb回调会在每次创建复制槽或请求其流式传输变更时调用,无论有多少变更已准备好输出。

typedef void (*LogicalDecodeStartupCB) (struct LogicalDecodingContext *ctx,
                                        OutputPluginOptions *options,
                                        bool is_init);

其中,is_init参数在创建复制槽时为 true,否则为 false。options指向一个输出插件可以设置的选项结构体:

typedef struct OutputPluginOptions
{
    OutputPluginOutputType output_type;
    bool        receive_rewrites;
} OutputPluginOptions;

output_type必须设为OUTPUT_PLUGIN_TEXTUAL_OUTPUTOUTPUT_PLUGIN_BINARY_OUTPUT。另见Section 49.6.3。如果receive_rewrites为 true,则在某些 DDL 操作期间由堆重写产生的变更,也会调用输出插件。这些变更对于处理 DDL 复制的插件有用,但需要特殊处理。

启动回调应当验证 ctx->output_plugin_options 中的选项。如果输出插 件需要保存状态,可以使用 ctx->output_plugin_private 来存储。

49.6.4.2. 关闭回调 #

当一个此前处于活动状态的复制槽不再使用时,就会调用可选的 shutdown_cb 回调。它可用于释放输出插件私有的资 源。此时未必是在删除该槽,也可能只是停止流式传输。

typedef void (*LogicalDecodeShutdownCB) (struct LogicalDecodingContext *ctx);

49.6.4.3. 事务开始回调 #

只要某个已提交事务的开始被解码,就会调用必需的 begin_cb 回调。已中止的事务及其内容永远不会被解 码。

typedef void (*LogicalDecodeBeginCB) (struct LogicalDecodingContext *ctx,
                                      ReorderBufferTXN *txn);

txn 参数包含该事务的元信息,例如它提交时的时间 戳以及它的 XID。

49.6.4.4. 事务结束回调 #

只要事务提交被解码,就会调用必需的 commit_cb 回 调。如果有被修改的行,那么在此之前,所有已修改行的 change_cb 回调都已经被调用过。

typedef void (*LogicalDecodeCommitCB) (struct LogicalDecodingContext *ctx,
                                       ReorderBufferTXN *txn,
                                       XLogRecPtr commit_lsn);

49.6.4.5. 更改回调 #

必需的change_cb回调会在事务中的每次行修改时被调用,不论该修改是INSERTUPDATEDELETE。即使原命令一次修改了多行,也会针对每一行单独调用此回调。

typedef void (*LogicalDecodeChangeCB) (struct LogicalDecodingContext *ctx,
                                       ReorderBufferTXN *txn,
                                       Relation relation,
                                       ReorderBufferChange *change);

参数ctxtxn所含的内容与begin_cbcommit_cb回调中的相同。此外,还会传入关系描述符relation(指向该行所属的关系),以及结构体change(描述该行的修改)。

Note

只有用户定义表中既不是不记录 WAL 的(见 UNLOGGED),也不是临时的(见 TEMPORARY or TEMP)更改,才能通过逻辑解码提 取出来。

49.6.4.6. 截断回调 #

可选的 truncate_cb 回调会在解码 TRUNCATE 命令时调用。

typedef void (*LogicalDecodeTruncateCB) (struct LogicalDecodingContext *ctx,
                                         ReorderBufferTXN *txn,
                                         int nrelations,
                                         Relation relations[],
                                         ReorderBufferChange *change);

这些参数与 change_cb 回调类似。不过,由于对通过外键 关联的表执行 TRUNCATE 时需要一起执行动作,所以该回调 接收的是关系数组,而不是单个关系。详见 TRUNCATE 语句的说明。

49.6.4.7. 源过滤回调 #

可选的 filter_by_origin_cb 回调用于判定,从 origin_id 重放而来的数据是否为输出插件所关心 的数据。

typedef bool (*LogicalDecodeFilterByOriginCB) (struct LogicalDecodingContext *ctx,
                                               RepOriginId origin_id);

ctx 参数的内容与其他回调相同。除了源本身之外,没有 其他信息可用。如果要表明来自传入节点的更改并不相关,则返回 true,这会 使这些更改被过滤掉;否则返回 false。对于被过滤掉的事务和更改,其他回调 都不会被调用。

在实现级联复制或多向复制方案时,这个回调很有用。按源过滤可以避免在这类 配置中同一更改被来回复制。虽然事务和更改本身也带有源信息,但通过这个 回调来过滤会明显更高效。

49.6.4.8. 通用消息回调 #

可选的message_cb回调会在每次解码出一条逻辑解码消息时被调用。

typedef void (*LogicalDecodeMessageCB) (struct LogicalDecodingContext *ctx,
                                        ReorderBufferTXN *txn,
                                        XLogRecPtr message_lsn,
                                        bool transactional,
                                        const char *prefix,
                                        Size message_size,
                                        const char *message);

其中,txn参数包含事务的元信息,例如提交时间戳及其 XID。但请注意,如果消息是非事务性的,并且记录该消息的事务尚未分配 XID,这个参数可能为 NULL。lsn包含消息的 WAL 位置。transactional表示消息是否以事务性方式发送。prefix是任意的、以空字符终止的前缀,可用于识别当前插件关注的消息。最后,message参数保存实际消息,其大小为message_size

应特别注意,确保输出插件视为有意义的消息前缀具有唯一性。使用扩展名或输 出插件自身的名称通常是一个不错的选择。

49.6.5. 产生输出的函数 #

为了真正产生输出,输出插件可以在 StringInfo 输出缓冲区 ctx->out 中写入数据,此时位于 begin_cbcommit_cbchange_cb 回调内部。 写入输出缓冲区之前,必须先调 用 OutputPluginPrepareWrite(ctx, last_write);写完 缓冲区之后,必须调用 OutputPluginWrite(ctx, last_write) 来执行写出。 last_write 指示某次写出是否为该回调的最后一次写 出。

下面的示例展示了如何把数据输出给输出插件的消费者:

OutputPluginPrepareWrite(ctx, true);
appendStringInfo(ctx->out, "BEGIN %u", txn->xid);
OutputPluginWrite(ctx, true);