PostgreSQL 源码树中的 contrib/test_decoding 子目录里有一个输出插件示例。
加载输出插件的方式是动态加载共享库,并以输出插件的名称作为库的基本名称。定位该库时使用常规库搜索路径。为了提供所需的输出插件回调,并表明这个库确实是输出插件,库必须提供一个名为_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操作,那么该操作会被忽略。
为了对变更进行解码、格式化和输出,输出插件可以使用后端的大部分常规基础设施,包括调用输出函数。可以对关系进行只读访问,但仅限于以下两种关系:由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()。
输出插件回调几乎可以用任意格式向消费者传递数据。对于某些用例,例如通过 SQL 查看更改,把数据返回为能够容纳任意数据的数据类型(例如 bytea)会很笨拙。如果输出插件只输出服务器编码中的文本数 据,它可以把 OutputPluginOptions.output_type 设置为 OUTPUT_PLUGIN_TEXTUAL_OUTPUT 而不是 OUTPUT_PLUGIN_BINARY_OUTPUT,并在 启动回调 中声明这一点。在这种情况下,所有数据都必须采用服务器编码,这样才能装入 text datum。断言开启的构建会检查这一点。
输出插件通过它所提供的各类回调获知正在发生的更改。
并发事务按提交顺序解码,并且只有属于某个特定事务的更改,才会在 begin 和 commit 回调之间被解码。 显式或隐式回滚的事务永远不会被解码。成功的保存点会按照它们在该事务中执行的顺序,被折叠进包含它们的事务中。
只有已经安全刷写到磁盘的事务才会被解码。这可能导致 COMMIT 在紧随其后的 pg_logical_slot_get_changes() 调用中不会立即被解 码,当 synchronous_commit 被设置为 off 时尤其如此。
可选的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_OUTPUT或OUTPUT_PLUGIN_BINARY_OUTPUT。另见Section 48.6.3。如果receive_rewrites为 true,则在某些 DDL 操作期间由堆重写产生的变更,也会调用输出插件。这些变更对于处理 DDL 复制的插件有用,但需要特殊处理。
启动回调应当验证 ctx->output_plugin_options 中的选项。如果输出插 件需要保存状态,可以使用 ctx->output_plugin_private 来存储。
当一个此前处于活动状态的复制槽不再使用时,就会调用可选的 shutdown_cb 回调。它可用于释放输出插件私有的资 源。此时未必是在删除该槽,也可能只是停止流式传输。
typedef void (*LogicalDecodeShutdownCB) (struct LogicalDecodingContext *ctx);
只要某个已提交事务的开始被解码,就会调用必需的 begin_cb 回调。已中止的事务及其内容永远不会被解 码。
typedef void (*LogicalDecodeBeginCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn);
txn 参数包含该事务的元信息,例如它提交时的时间 戳以及它的 XID。
只要事务提交被解码,就会调用必需的 commit_cb 回 调。如果有被修改的行,那么在此之前,所有已修改行的 change_cb 回调都已经被调用过。
typedef void (*LogicalDecodeCommitCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn,
XLogRecPtr commit_lsn);
必需的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 or TEMP)更改,才能通过逻辑解码提 取出来。
可选的 truncate_cb 回调会在解码 TRUNCATE 命令时调用。
typedef void (*LogicalDecodeTruncateCB) (struct LogicalDecodingContext *ctx,
ReorderBufferTXN *txn,
int nrelations,
Relation relations[],
ReorderBufferChange *change);
这些参数与 change_cb 回调类似。不过,由于对通过外键 关联的表执行 TRUNCATE 时需要一起执行动作,所以该回调 接收的是关系数组,而不是单个关系。详见 TRUNCATE 语句的说明。
可选的 filter_by_origin_cb 回调用于判定,从 origin_id 重放而来的数据是否为输出插件所关心 的数据。
typedef bool (*LogicalDecodeFilterByOriginCB) (struct LogicalDecodingContext *ctx,
RepOriginId origin_id);
ctx 参数的内容与其他回调相同。除了源本身之外,没有 其他信息可用。如果要表明来自传入节点的更改并不相关,则返回 true,这会 使这些更改被过滤掉;否则返回 false。对于被过滤掉的事务和更改,其他回调 都不会被调用。
在实现级联复制或多向复制方案时,这个回调很有用。按源过滤可以避免在这类 配置中同一更改被来回复制。虽然事务和更改本身也带有源信息,但通过这个 回调来过滤会明显更高效。
可选的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。
应特别注意,确保输出插件视为有意义的消息前缀具有唯一性。使用扩展名或输 出插件自身的名称通常是一个不错的选择。
为了真正产生输出,输出插件可以在 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);