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

基于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订阅)

适合单实例部署的服务,通过在服务内部维护订阅者列表,实现新任务的实时推送。

实现步骤:

  1. 扩展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),
	}
}
  1. 修改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
}
  1. 修改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内置的通知机制实现跨实例的任务推送同步。

实现步骤:

  1. 在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();
  1. 在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()
		}
	}
}
  1. 服务启动时启动监听协程:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 19:30:36