基于Go、gRPC、Postgres的任务列表:插入数据时自动流式传输实现咨询
任务列表应用的gRPC流式推送实现问题
我正在用Go、gRPC和Postgres开发任务列表应用,想知道调用PostItem接口插入新数据后,怎么实现数据的自动流式传输?是否需要订阅Postgres?或者不用订阅/发布订阅模式也能实现?
ProtoBuf 协议定义
syntax = "proto3"; package tasklist; import "google/protobuf/empty.proto"; service TodoList { rpc GetTasks(google.protobuf.Empty) returns (stream GetTasksResponse) {} rpc PostItem(PostItemRequest) returns (PostTaskRequest) {} } message Task { int64 id = 1; string name = 2; } message GetTasksResponse { Task task = 1; } message PostTaskRequest { Task Task = 1; } message PostItemResponse { bool result = 1; }
Postgres 表结构
create table Task ( id integer not null PRIMARY KEY, name varchar(10) not null );
Go 代码实现
func (s *server) GetTasks(_ *empty.Empty, stream pb.TaskList_GetTasksServer) error { // 如何在调用`PostTask`更新数据库后立即流式传输数据? <- <- for _, r := range s.requests { // 流式传输数据 } } func (s *server) PostTask(ctx context.Context, r *pb.PostTaskRequest) (*pb.PostTaskResponse, error) { // 在此处更新Postgres数据 return &pb.PostItemResponse{Result: true}, nil }
解决方案
方案一:内存级发布订阅(无需依赖Postgres订阅)
适合单实例部署的服务,通过在服务内部维护订阅者列表,实现新任务的实时推送。
实现步骤:
- 扩展
server结构体,添加锁、订阅者流列表和新任务通道:
import "sync" type server struct { pb.UnimplementedTodoListServer mu sync.Mutex streams []pb.TaskList_GetTasksServer taskCh chan *pb.Task } // 初始化server时创建通道 func NewServer() *server { return &server{ taskCh: make(chan *pb.Task, 100), } }
- 修改
GetTasks方法,注册当前流并监听新任务通道:
func (s *server) GetTasks(_ *empty.Empty, stream pb.TaskList_GetTasksServer) error { // 将当前流加入订阅列表 s.mu.Lock() s.streams = append(s.streams, stream) s.mu.Unlock() // 循环监听新任务,有数据就推送给客户端 for task := range s.taskCh { if err := stream.Send(&pb.GetTasksResponse{Task: task}); err != nil { // 移除失效的流 s.mu.Lock() for i, st := range s.streams { if st == stream { s.streams = append(s.streams[:i], s.streams[i+1:]...) break } } s.mu.Unlock() return err } } return nil }
- 修改
PostTask方法,插入数据库后推送新任务到通道:
func (s *server) PostTask(ctx context.Context, r *pb.PostTaskRequest) (*pb.PostItemResponse, error) { // 执行Postgres插入逻辑(示例) task := r.GetTask() _, err := s.db.ExecContext(ctx, "INSERT INTO Task(id, name) VALUES($1, $2)", task.Id, task.Name) if err != nil { return &pb.PostItemResponse{Result: false}, err } // 推送新任务给所有订阅者 s.taskCh <- task return &pb.PostItemResponse{Result: true}, nil }
优缺点:
- 优点:实现简单,无额外依赖,延迟低
- 缺点:仅支持单实例,服务重启后订阅会中断;多实例需结合Redis等中间件做跨实例消息同步
方案二:基于Postgres LISTEN/NOTIFY的订阅模式
适合多实例部署的场景,利用Postgres内置的通知机制实现跨实例的任务推送同步。
实现步骤:
- 在Postgres中创建触发器,插入任务后发送通知:
-- 创建通知函数 CREATE OR REPLACE FUNCTION notify_task_insert() RETURNS TRIGGER AS $$ BEGIN PERFORM pg_notify('task_insert', NEW.id || ',' || NEW.name); RETURN NEW; END; $$ LANGUAGE plpgsql; -- 绑定触发器到Task表 CREATE TRIGGER task_insert_trigger AFTER INSERT ON Task FOR EACH ROW EXECUTE FUNCTION notify_task_insert();
- 在Go服务中启动Postgres监听协程:
import ( "strings" "strconv" "github.com/lib/pq" ) func (s *server) startPostgresListener(ctx context.Context) error { // 建立Postgres连接 listener := pq.NewListener(s.dsn, 10*time.Second, time.Minute, func(n *pq.Notification) {}) if err := listener.Listen("task_insert"); err != nil { return err } defer listener.Close() for { select { case <-ctx.Done(): return nil case n := <-listener.Notify: if n == nil { continue } // 解析通知内容 parts := strings.Split(n.Extra, ",") id, _ := strconv.ParseInt(parts[0], 10, 64) newTask := &pb.Task{Id: id, Name: parts[1]} // 推送给所有订阅的gRPC流 s.mu.Lock() for _, stream := range s.streams { _ = stream.Send(&pb.GetTasksResponse{Task: newTask}) } s.mu.Unlock() } } }
- 服务启动时启动监听协程:
func main() { // 初始化server和数据库连接... go func() { if err := s.startPostgresListener(context.Background()); err != nil { log.Fatalf("Postgres listener failed: %v", err) } }() // 启动gRPC服务... }
优缺点:
- 优点:天然支持多实例同步,服务重启后可重新恢复监听
- 缺点:依赖Postgres特性,实现稍复杂;通知内容长度有限制(最多8KB)
总结
- 单实例服务优先选方案一,成本低易维护
- 多实例部署选方案二,或结合Redis Pub/Sub改造方案一实现跨实例同步
内容的提问来源于stack exchange,提问作者Pytan
相关产品推荐
相关产品推荐

