在Golang中实现向所有TCP连接发送数据的技术问询
实现Go TCP服务器的广播功能
要实现向所有当前连接发送数据,核心是维护一个全局的连接集合,并保证并发访问的安全性。以下是具体的改造方案:
关键修改点
- 给
Server结构体添加连接集合与互斥锁,用于安全管理所有活跃连接 - 在用户登录后将连接加入集合,连接关闭时从集合移除
- 实现广播方法,遍历所有连接并发送数据
- 处理消息时调用广播方法,将消息推送给所有在线用户
修改后的完整代码
package main import ( "fmt" "net" "strings" "sync" "time" ) type User struct { username string address string } type Message struct { from string payload []byte } type Server struct { listenAddr string ln net.Listener quitChan chan struct{} msgch chan Message conns map[net.Conn]bool // 存储所有活跃连接 mu sync.RWMutex // 保护conns的并发访问 } var user_db = []User{} func FindUserByAddress(array []User, address string) (username string) { for _, value := range array { if value.address == address { return value.username } } return "NOT_FOUND_IN_USER_DB" } func NewServer(listenAddr string) *Server { return &Server{ listenAddr: listenAddr, quitChan: make(chan struct{}), msgch: make(chan Message, 10), conns: make(map[net.Conn]bool), // 初始化连接集合 } } func (s *Server) Start() error { ln, err := net.Listen("tcp", s.listenAddr) if err != nil { return err } defer ln.Close() s.ln = ln go s.Acceptloop() go s.HandleMessages() // 单独的消息处理goroutine <-s.quitChan close(s.msgch) // 关闭所有活跃连接 s.mu.Lock() for conn := range s.conns { conn.Close() } s.mu.Unlock() return nil } func (s *Server) Acceptloop() { for { conn, err := s.ln.Accept() if err != nil { Log(err.Error()) continue } conn.Write([]byte("Welcome! Please, enter a username: ")) buffer := make([]byte, 128) n, err := conn.Read(buffer) if err != nil { Log(err.Error()) conn.Close() continue } username := strings.TrimSpace(string(buffer[:n])) conn.Write([]byte(fmt.Sprintf("Logged in as %v\n", username))) Log("New connection established from ", fmt.Sprintf("%v", conn.RemoteAddr()), ".\n Entered username: ", username) user_db = append(user_db, User{ username: username, address: conn.RemoteAddr().String(), }) // 将连接加入集合(加写锁) s.mu.Lock() s.conns[conn] = true s.mu.Unlock() go s.Read(conn) } } func (s *Server) Read(conn net.Conn) { defer func() { conn.Close() // 连接关闭时从集合移除(加写锁) s.mu.Lock() delete(s.conns, conn) s.mu.Unlock() Log("Connection closed: ", conn.RemoteAddr().String()) }() buffer := make([]byte, 2048) for { n, err := conn.Read(buffer) if err != nil { Log(err.Error()) return } s.msgch <- Message{ from: conn.RemoteAddr().String(), payload: buffer[:n], } conn.Write([]byte("Message delivered.\n")) } } // Broadcast 向所有活跃连接发送消息 func (s *Server) Broadcast(data []byte) { // 加读锁,因为只是遍历连接,不修改集合 s.mu.RLock() defer s.mu.RUnlock() for conn := range s.conns { // 逐个向连接写入数据,忽略写入错误(比如连接已断开) _, err := conn.Write(data) if err != nil { Log("Failed to send to ", conn.RemoteAddr().String(), ": ", err.Error()) } } } func (s *Server) HandleMessages() { for msg := range s.msgch { senderName := FindUserByAddress(user_db, msg.from) msgStr := strings.TrimSpace(string(msg.payload)) logMsg := fmt.Sprintf("[%s] %s: %s\n", msg.from, senderName, msgStr) Log(logMsg) // 构造广播消息并发送给所有用户 broadcastMsg := []byte(fmt.Sprintf("[%s] %s: %s\n", senderName, msg.from, msgStr)) s.Broadcast(broadcastMsg) } } func Log(args ...string) { fmt.Print("[", time.Now().Format("15:04:05"), "] ") for _, arg := range args { fmt.Print(arg) } fmt.Print("\n") } func main() { server := NewServer(":3000") Log("Starting server on port 3000...") if err := server.Start(); err != nil { Log(err.Error()) } }
核心细节解释
- 连接集合与并发安全:使用
sync.RWMutex保护connsmap,广播时是读操作(遍历所有连接),添加/删除连接是写操作,读写锁比互斥锁性能更好。 - 连接生命周期管理:在用户登录成功后将连接加入集合,
Read函数的defer块中关闭连接并从集合移除,确保集合中只保留活跃连接。 - 广播实现:
Broadcast方法遍历所有连接,逐个写入数据,忽略写入错误(避免单个连接故障影响整体广播)。 - 消息处理解耦:单独用
HandleMessagesgoroutine处理消息队列,收到消息后构造广播内容并调用Broadcast,逻辑更清晰。
内容的提问来源于stack exchange,提问作者SuperCarrot
相关产品推荐
相关产品推荐

