PostgreSQL的逻辑解码(Logical Decoding)功能允许开发者把数据库的物理变更日志(WAL)转换成结构化的、可读的逻辑变更流,这在数据同步、缓存刷新、审计、CDC等场景中应用极广。而实现这一能力的核心组件之一就是输出插件(Output Plugin)。PostgreSQL自带了一个名为pg_output的内置输出插件,它是逻辑复制发布订阅体系的默认载体,同时也是一个非常好的学习范例。很多开发者分不清pg_output和自定义解码插件的关系,也不清楚应该如何基于这套接口写出自己的插件,这篇文章就来彻底讲明白。

逻辑解码的整体架构与pg_output的位置
要理解pg_output,先要理解逻辑解码的三层结构。最底层是WAL(Write-Ahead Log),它记录的是物理层面的页级变更,比如某个数据页的某个偏移位置写入了什么内容,本身并不关心业务语义。中间层是Reorder Buffer(重排缓冲区),负责把物理WAL记录按照事务重新组织、按照提交顺序排序,并在事务提交时把该事务内的所有变更依次回放给上层。最上层才是输出插件,它接收已经整理好的逻辑变更事件,然后按照自己的协议格式把这些事件序列化后发送给消费端。
pg_output就是PostgreSQL官方实现的输出插件,源码位于PostgreSQL源码树的src/backend/replication/pgoutput/目录下。它的定位很明确:服务于逻辑复制的发布订阅(PUBLICATION/SUBSCRIPTION)机制。逻辑复制两端通过pg_output产生的二进制消息格式进行通信,这套格式是PostgreSQL内部定义的,并非人人可读的JSON或SQL文本。
值得注意的是,pg_output并不像test_decoding那样可以直接通过SQL创建槽位来体验,它需要配合逻辑复制协议使用。你可以在创建逻辑槽时显式指定它:
-- 这个语句可以执行,但pg_output的消息是二进制格式
-- 需要按逻辑复制协议解析才能读懂
SELECT * FROM pg_create_logical_replication_slot(
'my_slot', 'pgoutput');
真正日常调试时,更常用的做法是用test_decoding插件观察逻辑解码的输出,或者干脆自己写一个输出插件,把变更转换成JSON、Avro或者直接投递到消息队列。这也正是PostgreSQL把输出插件接口完全开放出来的意义所在:解码能力是通用的,输出格式和投递目标由你决定。
输出插件的回调接口详解
编写自定义输出插件,本质上是实现src/include/replication/output_plugin.h中定义的一组回调函数,并在插件加载时通过_PG_output_plugin_init入口注册它们。这个入口函数是PostgreSQL加载动态库时的约定入口,名字固定,不能改动。
核心回调包括以下几类:startup_cb在解码会话开始时被调用,可以在这里解析参数、初始化状态;shutdown_cb在会话结束时清理资源;begin_cb在一个事务开始回放时触发;commit_cb在事务提交时触发;change_cb是最重要的一个,每一行级别的INSERT、UPDATE、DELETE变更都会回调到这里;message_cb处理pg_logical_emit_message发出的自定义消息;truncate_cb处理TRUNCATE事件。此外还有针对两阶段提交的prepare_cb系列回调,如果不需要支持PREPARE TRANSACTION,可以留空。
一个最小的插件骨架如下:
#include "postgres.h"
#include "replication/logical.h"
#include "replication/output_plugin.h"
/* PG_MODULE_MAGIC 必须存在,否则数据库拒绝加载 */
PG_MODULE_MAGIC;
void _PG_output_plugin_init(OutputPluginCallbacks *cb);
void
_PG_output_plugin_init(OutputPluginCallbacks *cb)
{
/* 初始化回调结构,防止未设置的指针为随机值 */
memset(cb, 0, sizeof(OutputPluginCallbacks));
cb->startup_cb = pg_decode_startup;
cb->shutdown_cb = pg_decode_shutdown;
cb->begin_cb = pg_decode_begin_txn;
cb->commit_cb = pg_decode_commit_txn;
cb->change_cb = pg_decode_change;
}
void
pg_decode_change(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
Relation relation, ReorderBufferChange *change)
{
OutputPluginPrepareWrite(ctx, true);
switch (change->action)
{
case REORDER_BUFFER_CHANGE_INSERT:
appendBinaryStringInfo(ctx->out, "INSERT", 6);
break;
case REORDER_BUFFER_CHANGE_UPDATE:
appendBinaryStringInfo(ctx->out, "UPDATE", 6);
break;
case REORDER_BUFFER_CHANGE_DELETE:
appendBinaryStringInfo(ctx->out, "DELETE", 6);
break;
default:
appendBinaryStringInfo(ctx->out, "OTHER", 5);
}
OutputPluginWrite(ctx, true);
}
这里有一个新手最容易踩的坑:OutputPluginPrepareWrite和OutputPluginWrite必须成对调用。前者的第二个参数表示是否请求写端发送keepalive,后者的参数表示是否要求对端立即确认。如果你在change回调中积累了数据但始终不调用Write,数据会一直滞留在缓冲区,消费端什么都收不到。另外,写出去的内容最终进入ctx->out这个StringInfo缓冲区,你可以用appendStringInfo系列函数自由构造任意格式。
在change_cb里获取行数据也有一套固定套路。对于INSERT,新元组在change->data.tp.newtuple中;对于DELETE,旧元组在change->;data.tp.oldtuple中(注意只有在该表配置了REPLICA IDENTITY FULL或者主键时旧元组才可用);对于UPDATE,新旧元组都可能存在。元组需要结合关系描述符逐列解码:
HeapTuple tup = &change->data.tp.newtuple->tuple;
TupleDesc desc = RelationGetDescr(relation);
Datum values[MaxTupleAttributeNumber];
bool nulls[MaxTupleAttributeNumber];
heap_deform_tuple(tup, desc, values, nulls);
for (int i = 0; i < desc->natts; i++)
{
Form_pg_attribute attr = TupleDescAttr(desc, i);
if (nulls[i])
continue; /* NULL 列的处理 */
/* 根据 attr->atttypid 判断类型并格式化输出 */
}
这里强烈建议处理系统列以外的所有列时都通过类型OID判断,而不是假设列一定是text或int。真实业务表的类型千奇百怪,直接把Datum当字符串输出会导致乱码甚至崩溃。对于复杂类型,可以借助OidOutputFunctionCall配合getTypeOutputInfo获取该类型的输出函数,这是最稳妥的通用做法。
插件的编译、部署与验证
插件本质上是一个共享库动态链接文件(Linux下是.so,Windows下是.dll),编译时必须使用与数据库服务器完全一致的PostgreSQL开发头文件,且编译环境的major版本要与目标数据库一致,小版本不一致通常可以容忍,但跨大版本几乎必然出问题,因为回调结构体的定义在不同版本间有过调整。
编译命令参考如下(假设使用了PG_CONFIG环境变量指向目标安装):
gcc -fPIC -c my_decoder.c -o my_decoder.o \
$(pg_config --cflags)
gcc -shared my_decoder.o -o my_decoder.so
cp my_decoder.so $(pg_config --pkglibdir)
部署完成后还要调整两个关键配置:在postgresql.conf中设置wal_level = logical,同时把max_replication_slots调到大于0(默认10已经够用)。修改后需要重启数据库。别忘了pg_hba.conf中如果打算让远程客户端消费,还需要加入支持复制的认证条目,认证方法可以选scram-sha-256。
验证阶段可以先创建槽位直接消费:
SELECT * FROM pg_create_logical_replication_slot(
'demo_slot', 'my_decoder');
-- 制造一点变更
CREATE TABLE t1(id int primary key, name text);
INSERT INTO t1 VALUES (1, 'hello');
-- 拉取并消费解码结果
SELECT lsn, xid, data FROM pg_logical_slot_get_changes(
'demo_slot', NULL, NULL);
如果一切正常,你会看到刚才那条INSERT被插件转换成了自定义格式的输出。调试阶段建议优先使用pg_logical_slot_peek_changes,它只查看不推进槽位确认位点,方便反复测试同一段WAL而不用反复制造数据。这一点在开发回调逻辑时能节省大量时间。
最后说几个常见的坑。第一,逻辑解码依赖REPLICA IDENTITY,没有主键的表在UPDATE或DELETE时旧元组拿不到,插件会收到空指针,必须判空处理。第二,长事务会撑大Reorder Buffer,解码滞后可能把disk上的临时文件越堆越多,生产环境要监控pg_replication_slots视图中的confirmed_flush_lsn与当前LSN的差距。第三,槽位是持久的,测试完记得用pg_drop_replication_slot清理,否则WAL无法回收,磁盘迟早被撑满。
掌握输出插件这套接口之后,你不仅能写出满足自己格式需求的解码器,也能更透彻地理解pg_output乃至Debezium、pg_stat_stream_replication这些生态工具的工作原理。逻辑解码是PostgreSQL最具扩展性的能力之一,值得每个做数据基础设施的工程师深入研究。
PostgreSQLpg_output插件逻辑解码修改时间:2026-09-06 00:39:00