You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 05:09:20