Golang使用lib/pq实现PostgreSQL触发器及修改行返回示例
Golang 基于lib/pq驱动使用PostgreSQL触发器获取修改后行数据实现
你写的trigger_set_timestamp是标准的行级BEFORE触发器实现,本身在执行INSERT/UPDATE操作时会自动修改updated_at字段,最终写入表的就是修改完成的NEW行,不需要额外依赖通知框架就能拿到结果;如果需要异步感知表变更,再配合PostgreSQL的LISTEN/NOTIFY机制推送数据即可。
一、PostgreSQL端触发器配置
首先完成基础表和触发器的创建,以下是可直接运行的SQL:
-- 建演示用表 CREATE TABLE IF NOT EXISTS demo_records ( id SERIAL PRIMARY KEY, content TEXT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ); -- 自动更新时间戳的触发器函数 CREATE OR REPLACE FUNCTION trigger_set_timestamp() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at = NOW(); RETURN NEW; END; $$ LANGUAGE plpgsql; -- 绑定BEFORE行级触发器,插入、更新时自动触发 CREATE TRIGGER set_demo_timestamp BEFORE INSERT OR UPDATE ON demo_records FOR EACH ROW EXECUTE FUNCTION trigger_set_timestamp();
如果需要异步推送变更数据到Go服务,新增支持NOTIFY的触发器函数即可,触发器会在修改完字段后把整行数据作为通知载荷发送:
CREATE OR REPLACE FUNCTION trigger_notify_change() RETURNS TRIGGER AS $$ DECLARE payload JSON; BEGIN -- 执行字段修改逻辑 NEW.updated_at = NOW(); -- 将修改后的完整行转为JSON作为通知内容 payload = row_to_json(NEW); -- 发送通知到指定通道 PERFORM pg_notify('record_change_channel', payload::text); RETURN NEW; END; $$ LANGUAGE plpgsql; -- 绑定异步通知触发器 CREATE TRIGGER notify_demo_change BEFORE INSERT OR UPDATE ON demo_records FOR EACH ROW EXECUTE FUNCTION trigger_notify_change();
二、Golang端lib/pq实现代码
分两种常用场景实现,分别对应同步拿返回值、异步接通知的需求。
场景1:同步执行DML直接获取修改后行数据
BEFORE触发器修改NEW并RETURN NEW后,最终写入的值可以通过RETURNING子句直接拿到,不需要额外查询:
package main import ( "database/sql" "fmt" "time" "github.com/lib/pq" ) type DemoRecord struct { ID int `db:"id"` Content string `db:"content"` CreatedAt time.Time `db:"created_at"` UpdatedAt time.Time `db:"updated_at"` } func main() { connStr := "host=localhost port=5432 user=postgres password=your_db_pass dbname=test_db sslmode=disable" db, err := sql.Open("postgres", connStr) if err != nil { panic(err) } defer db.Close() // 插入数据,直接RETURNING拿到触发器自动设置的updated_at var newRecord DemoRecord err = db.QueryRow( `INSERT INTO demo_records (content) VALUES ($1) RETURNING id, content, created_at, updated_at`, "测试内容", ).Scan(&newRecord.ID, &newRecord.Content, &newRecord.CreatedAt, &newRecord.UpdatedAt) if err != nil { panic(err) } fmt.Printf("插入后返回的行(触发器已自动设置updated_at):%+v\n", newRecord) // 更新数据,同样通过RETURNING拿到触发器刷新后的行 var updatedRecord DemoRecord err = db.QueryRow( `UPDATE demo_records SET content = $1 WHERE id = $2 RETURNING id, content, created_at, updated_at`, "修改后的测试内容", newRecord.ID, ).Scan(&updatedRecord.ID, &updatedRecord.Content, &updatedRecord.CreatedAt, &updatedRecord.UpdatedAt) if err != nil { panic(err) } fmt.Printf("更新后返回的行(触发器已自动刷新updated_at):%+v\n", updatedRecord) // 需要异步监听的话启动监听协程 go listenChange(connStr) select {} }
场景2:异步监听通知获取变更行
用lib/pq自带的Listener监听通道,接收触发器推送的变更数据:
func listenChange(connStr string) { // 建立监听专用连接 listener := pq.NewListener(connStr, 10*time.Second, time.Minute, func(ev pq.ListenerEventType, err error) { if err != nil { fmt.Printf("监听连接异常:%v\n", err) } }) err := listener.Listen("record_change_channel") if err != nil { panic(err) } fmt.Println("开始监听record_change_channel通道变更...") for { select { case n := <-listener.Notify: // n.Extra 就是触发器发送的JSON格式修改后行数据,可直接反序列化使用 fmt.Printf("收到变更行数据:%s\n", n.Extra) case <-time.After(90 * time.Second): // 定时心跳保活 go func() { _ = listener.Ping() }() } } }
注意事项
- 只有
BEFORE行级触发器对NEW字段的修改会被持久化,AFTER触发器中修改NEW不会影响最终写入的数据 - 同步场景下必须在DML语句末尾加
RETURNING子句指定返回字段,才能直接拿到触发器修改后的完整行,不需要额外执行SELECT查询 - NOTIFY机制的消息载荷上限为8000字节,大字段行不要直接序列化整行发送,推荐只传递变更行主键,收到通知后再查表拿最新数据,避免消息截断
- LISTEN/NOTIFY不保证消息可靠投递,监听连接断开期间产生的变更不会在重连后补发,强一致性场景优先用同步RETURNING的方式拿数据
内容的提问来源于stack exchange,提问作者Ajoy Das
相关产品推荐
相关产品推荐

