48.6. 逻辑解码输出插件 #
PostgreSQL 源码树中的
contrib/test_decoding
子目录里有一个输出插件示例。
48.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_cb,change_cb 和 commit_cb 回调是必需的,而 startup_cb,filter_by_origin_cb,truncate_cb,以及 shutdown_cb 是可选的。如果没有设置 truncate_cb,但要解码一个 TRUNCATE 操作,那么该操作会被忽略。
48.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()。
48.6.3. 输出模式 #
输出插件回调几乎可以用任意格式向消费者传递数据。对于某些用例,例如通过
SQL 查看更改,把数据返回为能够容纳任意数据的数据类型(例如
bytea)会很不方便。如果输出插件只输出采用服务器编码的文本数据,它可以在启动回调中,将 OutputPluginOptions.output_type 设置为
OUTPUT_PLUGIN_TEXTUAL_OUTPUT 而不是
OUTPUT_PLUGIN_BINARY_OUTPUT,以声明这一点。在这种情况下,所有数据都必须采用服务器编码,这样才能装入 text
datum。断言开启的构建会检查这一点。
48.6.4. 输出插件回调 #
输出插件通过它所提供的各类回调获知正在发生的更改。
并发事务按提交顺序解码,并且只有属于某个特定事务的更改,才会在
begin 和 commit 回调之间被解码。显式或隐式回滚的事务永远不会被解码。成功的保存点会按照它们在该事务中执行的顺序,被折叠进包含它们的事务中。
注意
只有已经安全刷写到磁盘的事务才会被解码。因此,当
synchronous_commit 设置为 off 时,COMMIT 可能不会在紧随其后的
pg_logical_slot_get_changes() 调用中立即被解码。
48.6.4.1. 启动回调 #
只要创建复制槽,或者要求某个复制槽开始流式传输更改,就会调用可选的
startup_cb 回调,而不管当前是否已经有准备好输出的更改。
typedef void (*LogicalDecodeStartupCB) (struct LogicalDecodingContext *ctx,
OutputPluginOptions *options,
bool is_init);
当复制槽正在创建时,is_init 参数为真,否则为假。options 指向一个结构体,输出插件可以在其中设置选项:
typedef struct OutputPluginOptions
{
OutputPluginOutputType output_type;
bool receive_rewrites;
} OutputPluginOptions;
output_type 必须设置为
OUTPUT_PLUGIN_TEXTUAL_OUTPUT 或
OUTPUT_PLUGIN_BINARY_OUTPUT。另见第 48.6.3 节。如果
receive_rewrites 为真,则在某些 DDL 操作期间由堆重写产生的更改也会传给输出插件。这对处理 DDL 复制的插件很有用,但需要特殊处理。
启动回调应当验证
ctx->output_plugin_options 中的选项。如果输出插件需要保存状态,可以使用
ctx->output_plugin_private 来存储。
48.6.4.2. 关闭回调 #
当一个此前处于活动状态的复制槽不再使用时,就会调用可选的
shutdown_cb 回调。它可用于释放输出插件私有的资源。此时只是停止流式传输,并不一定删除该槽。
typedef void (*LogicalDecodeShutdownCB) (struct LogicalDecodingContext *ctx);
48.6.4.3. 事务开始回调 #
只要某个已提交事务的开始被解码,就会调用必需的
begin_cb 回调。已中止的事务及其内容永远不会被解码。
typedef void (*LogicalDecodeBeginCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn);
txn 参数包含该事务的元信息,例如它提交时的时间戳以及它的 XID。
48.6.4.4. 事务结束回调 #
只要事务提交被解码,就会调用必需的 commit_cb 回调。如果有被修改的行,那么在此之前,所有已修改行的
change_cb 回调都已经被调用过。
typedef void (*LogicalDecodeCommitCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn,
XLogRecPtr commit_lsn);
48.6.4.5. 更改回调 #
必需的 change_cb 回调会在事务中的每次行修改时被调用,不论该修改是 INSERT,UPDATE 或 DELETE。即使原命令一次修改了多行,也会针对每一行单独调用此回调。
typedef void (*LogicalDecodeChangeCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn,
Relation relation,
ReorderBufferChange *change);
参数 ctx 和 txn 所含的内容与 begin_cb 和 commit_cb 回调中的相同。此外,还会传入关系描述符 relation(指向该行所属的关系),以及结构体 change(描述该行的修改)。
注意
只有用户定义表中的更改才能通过逻辑解码提取,而且这些表不能是不记录 WAL 的表(见 UNLOGGED),也不能是临时表(见 TEMPORARY 或 TEMP)。
48.6.4.6. 截断回调 #
可选的 truncate_cb 回调会在解码
TRUNCATE 命令时调用。
typedef void (*LogicalDecodeTruncateCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn,
int nrelations,
Relation relations[],
ReorderBufferChange *change);
这些参数与 change_cb 回调类似。不过,由于对通过外键关联的表执行 TRUNCATE 时需要一起执行动作,所以该回调接收的是关系数组,而不是单个关系。详见 TRUNCATE 语句的说明。
48.6.4.7. 源过滤回调 #
可选的 filter_by_origin_cb 回调用于判定,从
origin_id 重放而来的数据是否为输出插件所关心的数据。
typedef bool (*LogicalDecodeFilterByOriginCB) (struct LogicalDecodingContext *ctx,
RepOriginId origin_id);
ctx 参数的内容与其他回调相同。除了源本身之外,没有其他信息可用。如果要表明来自传入节点的更改并不相关,则返回 true,这会使这些更改被过滤掉;否则返回 false。对于被过滤掉的事务和更改,其他回调都不会被调用。
在实现级联复制或多向复制方案时,这个回调很有用。按源过滤可以避免在这类配置中同一更改被来回复制。虽然事务和更改本身也带有源信息,但通过这个回调来过滤会明显更高效。
48.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。
应特别注意,确保输出插件视为有意义的消息前缀具有唯一性。使用扩展的名称或输出插件自身的名称通常是一个不错的选择。
48.6.5. 产生输出的函数 #
为了真正产生输出,输出插件可以在
StringInfo 输出缓冲区 ctx->out
中写入数据,此时位于 begin_cb、commit_cb 或 change_cb 回调内部。写入输出缓冲区之前,必须先调用 OutputPluginPrepareWrite(ctx, last_write);写完缓冲区之后,必须调用
OutputPluginWrite(ctx, last_write) 来执行写出。last_write 指示某次写出是否为该回调的最后一次写出。
下面的示例展示了如何把数据输出给输出插件的消费者:
OutputPluginPrepareWrite(ctx, true); appendStringInfo(ctx->out, "BEGIN %u", txn->xid); OutputPluginWrite(ctx, true);