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

在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())
	}
}

核心细节解释

  1. 连接集合与并发安全:使用sync.RWMutex保护conns map,广播时是读操作(遍历所有连接),添加/删除连接是写操作,读写锁比互斥锁性能更好。
  2. 连接生命周期管理:在用户登录成功后将连接加入集合,Read函数的defer块中关闭连接并从集合移除,确保集合中只保留活跃连接。
  3. 广播实现:Broadcast方法遍历所有连接,逐个写入数据,忽略写入错误(避免单个连接故障影响整体广播)。
  4. 消息处理解耦:单独用HandleMessages goroutine处理消息队列,收到消息后构造广播内容并调用Broadcast,逻辑更清晰。

内容的提问来源于stack exchange,提问作者SuperCarrot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 03:14:59