并发模型(CSP理论、Goroutine轻量级线程)#
基础概念#
并发: 多个任务在同一时间段内交替执行, 但不一定同时执行, 也就是任务之间的切换是有序的, 任务之间是有依赖关系的.
并行: 多个任务在同一时刻同时执行
CSP理论:
- 核心思想: 不要通过共享内存来通信,而要通过通信来共享内存
- 组成要素:
Process(进程\协程)、Channel(通道)
应用场景#
Goroutine:
- 高并发Web服务器:每个请求一个Goroutine
- 数据处理流水线:多个处理阶段并行执行
- 实时消息推送:WebSocket连接管理
- 定时任务调度:后台定时执行任务
- 并发爬虫:同时抓取多个网页
传统多线程与Goroutine的区别:
| 特性 | 传统线程 | Goroutine |
|---|---|---|
| 创建成本 | 1-2MB | 2KB |
| 创建速度 | 慢 | 快 |
| 调度方式 | 操作系统调度 | Go运行时调度 |
| 上下文切换 | 完整线程切换 | 用户态轻量切换 |
通道(Channel)#
Channel是Go语言中各个并发执行体间的通信机制,是类型相关的管道,用于在goroutine之间传递数据和同步执行。
不要通过共享内存来通信,而应通过通信来共享内存
// 创建channel
ch1 := make(chan int) // 无缓冲channel
ch2 := make(chan int, 10) // 有缓冲channel,容量10
// 操作
ch1 <- 42 // 发送数据到channel
value := <-ch1 // 从channel接收数据
close(ch1) // 关闭channel
// 特殊用法
value, ok := <-ch1 // 检查channel是否关闭go缓冲区区别:
- 无缓冲channel:同步通信,发送和接收必须同时准备好
- 有缓冲channel:异步通信,缓冲区满时发送阻塞,空时接收阻塞
package main
import (
"fmt"
"math/rand"
"time"
)
type Order struct {
ID int
UserID string
Amount float64
Status string
CreatedAt time.Time
}
func orderProducer(orderChan chan<- Order, numOrders int) {
for i := 1; i <= numOrders; i++ {
order := Order{
ID: i,
UserID: fmt.Sprintf("user%d", rand.Intn(100)),
Amount: rand.Float64() * 1000,
Status: "pending",
CreatedAt: time.Now(),
}
orderChan <- order
fmt.Printf("生成订单: ID=%d, 用户=%s, 金额=¥%.2f\n",
order.ID, order.UserID, order.Amount)
time.Sleep(time.Millisecond * 100) // 模拟生成间隔
}
close(orderChan)
}
func orderProcessor(orderChan <-chan Order, resultChan chan<- Order) {
for order := range orderChan {
// 模拟订单处理逻辑
processingTime := time.Duration(rand.Intn(300)) * time.Millisecond
time.Sleep(processingTime)
// 更新订单状态
if order.Amount > 500 {
order.Status = "verified" // 大额订单需要验证
} else {
order.Status = "completed"
}
resultChan <- order
}
time.Sleep(time.Second * 1)
close(resultChan)
}
func orderResultCollector(resultChan <-chan Order, done chan<- bool) {
processedCount := 0
for order := range resultChan {
processedCount++
fmt.Printf("处理完成: 订单ID=%d, 状态=%s, 金额=¥%.2f\n",
order.ID, order.Status, order.Amount)
}
fmt.Printf("所有订单处理完成! 总计: %d 个订单\n", processedCount)
done <- true
}
func main() {
rand.Seed(time.Now().UnixNano())
// 创建管道
orderChan := make(chan Order, 10)
resultChan := make(chan Order, 10)
done := make(chan bool)
// 启动订单生成器
go orderProducer(orderChan, 20)
// 启动多个订单处理器(工人)
for i := 1; i <= 3; i++ {
go orderProcessor(orderChan, resultChan)
}
// 启动结果收集器
go orderResultCollector(resultChan, done)
// 等待所有处理完成
<-done
}goselect多路复用:同时等待多个channel操作。
package main
import (
"fmt"
"time"
)
func main() {
// 1. select多路复用
ch1 := make(chan string)
ch2 := make(chan string)
go func() {
time.Sleep(1 * time.Second)
ch1 <- "来自ch1"
}()
go func() {
time.Sleep(2 * time.Second)
ch2 <- "来自ch2"
}()
for i := 0; i < 2; i++ {
select {
case msg1 := <-ch1:
fmt.Println("收到:", msg1)
case msg2 := <-ch2:
fmt.Println("收到:", msg2)
case <-time.After(3 * time.Second): // 超时控制
fmt.Println("超时!")
return
}
}
// 2. 定时器与Ticker
ticker := time.NewTicker(500 * time.Millisecond)
done := make(chan bool)
go func() {
for {
select {
case <-done:
return
case t := <-ticker.C:
fmt.Println("定时触发 at", t.Format("15:04:05"))
}
}
}()
time.Sleep(2 * time.Second)
ticker.Stop()
done <- true
fmt.Println("Ticker停止")
// 3. 工作池模式
jobs := make(chan int, 100)
results := make(chan int, 100)
// 启动3个worker
for w := 1; w <= 3; w++ {
go worker(w, jobs, results)
}
// 发送5个任务
for j := 1; j <= 5; j++ {
jobs <- j
}
close(jobs)
// 收集结果
for i:=0; i<=5; i++ {
value := <-results
fmt.Printf("Worker 处理结果 %d\n", value)
}
}
func worker(id int, jobs <-chan int, results chan<- int) {
for j := range jobs {
fmt.Printf("Worker %d 处理任务 %d\n", id, j)
time.Sleep(time.Second)
results <- j * 2
}
}go同步原语(sync.Mutex互斥锁、sync.WaitGroup等待组、sync.RWMutex读写锁)#
为什么需要同步原语? 在并发编程中,多个goroutine同时访问共享资源时会产生竞态条件,导致数据不一致。
sync.Mutex(互斥锁)
- 保证同一时间只有一个goroutine能访问共享资源
- 两个方法:
Lock()和Unlock() - 使用后必须释放,否则会导致死锁
sync.RWMutex(读写锁)
- 允许多个读操作或一个写操作
- 读锁:
RLock()/RUnlock()(共享锁) - 写锁:
Lock()/Unlock()(互斥锁) - 适合”读多写少”的场景
sync.WaitGroup(等待组)
- 等待一组goroutine完成工作
- 三个方法:
Add()、Done()、Wait() - 用于协调多个goroutine的执行顺序
var wg sync.WaitGroup
func main() {
for i := 0; i < 3; i++ {
wg.Add(1) // 计数器+1
go worker(i)
}
wg.Wait() // 等待所有goroutine完成
}gopackage main
import (
"fmt"
"sync"
"time"
)
type Inventory struct {
stock int // 库存数量
rwMutex sync.RWMutex // 保护库存
}
// 查询库存(使用读写锁的读锁)
func (inv *Inventory) getStock() int {
inv.rwMutex.RLock()
defer inv.rwMutex.RUnlock()
return inv.stock
}
// 扣减库存(使用互斥锁和读写锁的写锁)
func (inv *Inventory) deductStock(userID int, quantity int) bool {
// 先检查库存(读锁)
inv.rwMutex.Lock()
defer inv.rwMutex.Unlock()
// 检查,防止超卖
if inv.stock < quantity {
fmt.Printf("用户%d: 库存不足,扣减失败\n", userID)
return false
}
inv.stock -= quantity
fmt.Printf("用户%d: 成功购买%d件,剩余库存%d\n", userID, quantity, inv.stock)
return true
}
func main() {
inventory := &Inventory{stock: 10} // 初始库存10件
var wg sync.WaitGroup
// 模拟100个用户同时抢购
for i := 1; i <= 100; i++ {
wg.Add(1)
go func(userID int) {
defer wg.Done()
inventory.deductStock(userID, 1)
}(i)
}
wg.Wait()
fmt.Printf("最终库存: %d\n", inventory.getStock())
}go原子操作替代简单锁:
import "sync/atomic"
type Counter struct {
value int64
}
func (c *Counter) Increment() {
atomic.AddInt64(&c.value, 1) // 比mutex性能更好
}
func (c *Counter) Decrement() {
atomic.AddInt64(&c.value, -1)
}
func (c *Counter) Value() int64 {
return atomic.LoadInt64(&c.value)
}
func main() {
var counter Counter
var wg sync.WaitGroup
// 模拟100个用户同时点赞
for i := 1; i <= 100; i++ {
wg.Add(1)
go func() {
defer wg.Done()
counter.Increment()
}()
}
// 模拟10个用户取消点赞
for i := 1; i <= 10; i++ {
wg.Add(1)
go func() {
defer wg.Done()
counter.Decrement()
}()
}
wg.Wait()
}go并发模式(生产者-消费者、扇入/扇出、Pipeline)#
生产者-消费者模式
- 生产者-消费者模式是并发编程中最经典的模式之一,解决生产者和消费者速度不匹配的问题。
- 就像食堂阿姨打饭和同学们吃饭,赚钱和花钱的关系是一样的。
扇入/扇出模式:
- 扇出:一个channel分发给多个goroutine处理(一产多消)
- 扇入:多个channel合并到一个channel(多产一消)
Pipeline模式:将复杂任务分解为多个处理阶段,每个阶段通过channel连接,形成处理流水线。
package main
import (
"fmt"
"math/rand"
"sync"
"time"
)
// 订单结构
type Order struct {
ID int
UserID int
Amount float64
Status string
CreateAt time.Time
}
// 1. 生产者-消费者:订单生成与处理
func orderProducer(orderCh chan<- Order, count int) {
defer close(orderCh)
for i := 1; i <= count; i++ {
order := Order{
ID: i,
UserID: rand.Intn(1000) + 1,
Amount: rand.Float64() * 1000,
Status: "pending",
CreateAt: time.Now(),
}
fmt.Printf("生成订单: ID=%d, 金额=¥%.2f\n", order.ID, order.Amount)
orderCh <- order
time.Sleep(100 * time.Millisecond) // 模拟生成间隔
}
}
func orderConsumer(orderCh <-chan Order, wg *sync.WaitGroup) {
defer wg.Done()
for order := range orderCh {
// 模拟订单处理
time.Sleep(200 * time.Millisecond)
order.Status = "processed"
fmt.Printf("处理订单: ID=%d, 状态=%s\n", order.ID, order.Status)
}
}
// 2. 扇出模式:一个订单流分发给多个处理器
func fanOutProcessor(input <-chan Order, workerID int, wg *sync.WaitGroup) {
defer wg.Done()
for order := range input {
// 不同的处理器做不同的处理
switch workerID {
case 1:
// 处理器1:计算折扣
discount := order.Amount * 0.1
fmt.Printf("处理器%d计算折扣: 订单%d 优惠¥%.2f\n",
workerID, order.ID, discount)
case 2:
// 处理器2:发送通知
fmt.Printf("处理器%d发送通知: 订单%d创建成功\n",
workerID, order.ID)
case 3:
// 处理器3:记录日志
fmt.Printf("处理器%d记录日志: 订单%d金额¥%.2f\n",
workerID, order.ID, order.Amount)
}
time.Sleep(150 * time.Millisecond)
}
}
// 3. 扇入模式:多个数据源合并
func fanInProducer(producerID int, output chan<- Order) {
defer fmt.Printf("生产者%d结束\n", producerID)
for i := 1; i <= 3; i++ {
order := Order{
ID: producerID*100 + i,
UserID: producerID,
Amount: float64(producerID*100 + i),
Status: "new",
}
fmt.Printf("生产者%d生成订单: %d\n", producerID, order.ID)
output <- order
time.Sleep(time.Duration(producerID) * 100 * time.Millisecond)
}
}
// 4. Pipeline模式:订单处理流水线
func validationStage(input <-chan Order) <-chan Order {
output := make(chan Order, 10)
go func() {
defer close(output)
for order := range input {
// 第一阶段:订单验证
time.Sleep(50 * time.Millisecond)
if order.Amount > 0 {
order.Status = "validated"
fmt.Printf("验证通过: 订单%d\n", order.ID)
output <- order
} else {
fmt.Printf("验证失败: 订单%d金额异常\n", order.ID)
}
}
}()
return output
}
func paymentStage(input <-chan Order) <-chan Order {
output := make(chan Order, 10)
go func() {
defer close(output)
for order := range input {
// 第二阶段:支付处理
time.Sleep(100 * time.Millisecond)
order.Status = "paid"
fmt.Printf("支付成功: 订单%d\n", order.ID)
output <- order
}
}()
return output
}
func shippingStage(input <-chan Order) <-chan Order {
output := make(chan Order, 10)
go func() {
defer close(output)
for order := range input {
// 第三阶段:发货处理
time.Sleep(150 * time.Millisecond)
order.Status = "shipped"
fmt.Printf("发货完成: 订单%d\n", order.ID)
output <- order
}
}()
return output
}
func main() {
rand.Seed(time.Now().UnixNano())
fmt.Println("=== 1. 生产者-消费者模式演示 ===")
// 创建订单channel
orderCh := make(chan Order, 5)
var wg sync.WaitGroup
// 启动消费者
wg.Add(2)
go orderConsumer(orderCh, &wg)
go orderConsumer(orderCh, &wg)
// 启动生产者
go orderProducer(orderCh, 6)
wg.Wait()
fmt.Println("\n=== 2. 扇出模式演示 ===")
// 扇出:一个输入,多个处理器
fanOutCh := make(chan Order, 10)
var fanOutWg sync.WaitGroup
// 启动3个处理器
fanOutWg.Add(3)
for i := 1; i <= 3; i++ {
go fanOutProcessor(fanOutCh, i, &fanOutWg)
}
// 生产一些测试数据
go func() {
for i := 1; i <= 6; i++ {
fanOutCh <- Order{ID: i, Amount: float64(i * 100)}
}
close(fanOutCh)
}()
fanOutWg.Wait()
fmt.Println("\n=== 3. 扇入模式演示 ===")
// 扇入:多个生产者,一个输出
fanInCh := make(chan Order, 10)
// 启动3个生产者
for i := 1; i <= 3; i++ {
go fanInProducer(i, fanInCh)
}
// 收集结果
go func() {
time.Sleep(2 * time.Second)
close(fanInCh)
}()
fmt.Println("收集到的订单:")
for order := range fanInCh {
fmt.Printf(" 订单ID: %d, 金额: ¥%.2f\n", order.ID, order.Amount)
}
fmt.Println("\n=== 4. Pipeline模式演示 ===")
// 创建初始输入
pipelineInput := make(chan Order, 10)
// 构建流水线
validatedOrders := validationStage(pipelineInput)
paidOrders := paymentStage(validatedOrders)
shippedOrders := shippingStage(paidOrders)
// 发送测试订单到流水线
go func() {
for i := 1; i <= 3; i++ {
pipelineInput <- Order{
ID: i,
UserID: i * 10,
Amount: float64(i * 50),
Status: "new",
}
}
close(pipelineInput)
}()
// 收集最终结果
fmt.Println("流水线处理结果:")
for order := range shippedOrders {
fmt.Printf(" 完成: 订单%d, 状态: %s\n", order.ID, order.Status)
}
fmt.Println("\n🎉 所有并发模式演示完成!")
}go带错误处理的并发模式:
package main
import (
"errors"
"fmt"
"sync"
"time"
)
type Result struct {
Value int
Error error
}
// 带错误处理的生产者-消费者
func safeProducer(ch chan<- Result, wg *sync.WaitGroup) {
defer wg.Done()
for i := 1; i <= 5; i++ {
// 模拟偶尔出错
if i == 3 {
ch <- Result{Error: errors.New("模拟错误")}
} else {
ch <- Result{Value: i}
}
time.Sleep(100 * time.Millisecond)
}
close(ch)
}
func safeConsumer(ch <-chan Result, wg *sync.WaitGroup) {
defer wg.Done()
for result := range ch {
if result.Error != nil {
fmt.Printf("处理出错: %v\n", result.Error)
} else {
fmt.Printf("处理成功: %d\n", result.Value)
}
}
}
func main() {
fmt.Println("=== 带错误处理的并发 ===")
resultCh := make(chan Result, 5)
var safeWg sync.WaitGroup
safeWg.Add(2)
go safeProducer(resultCh, &safeWg)
go safeConsumer(resultCh, &safeWg)
safeWg.Wait()
}go网络基础(TCP/IP协议、三次握手、端口概念)#
核心概念#
TCP/IP协议族是互联网通信的基础,就像现实世界的邮政系统一样,负责数据的可靠传输。
TCP/IP四层模型(类比快递系统):
- 应用层 - 你的信件内容(HTTP、FTP、SMTP)
- 传输层 - 快递包装和物流单(TCP、UDP)
- 网络层 - 地址和路由(IP协议)
- 网络接口层 - 运输工具(以太网、WiFi)
TCP vs UDP 核心区别:
-
TCP - 可靠传输,像打电话
- 需要建立连接(三次握手)
- 保证数据顺序和完整性
- 自动重传丢失的数据
- 适合:网页浏览、文件传输、邮件
-
UDP - 不可靠传输,像发短信
- 无需建立连接
- 不保证数据到达顺序
- 可能丢失数据包
- 适合:视频流、游戏、DNS查询
三次握手详解#
TCP连接建立过程 - 就像确认双方都能正常沟通:

为什么需要三次握手?
- 防止历史连接:避免网络延迟导致的重复连接请求
- 同步序列号:确保双方都知道对方的起始序号
- 确认双向可达:证明客户端和服务器都能正常收发数据
端口概念与应用#
端口的作用 - 就像公司的分机号,员工号:
- 80: HTTP - 网页浏览
- 443: HTTPS - 安全网页
- 22: SSH - 安全远程登录
- 53: DNS - 域名解析
- 3306: MySQL - 数据库
- 5432: PostgreSQL - 数据库
- 6379: Redis - 缓存数据库
package main
import (
"fmt"
"net"
"sync"
)
type ChatRoom struct {
clients map[net.Conn]string
mutex sync.RWMutex
}
func NewChatRoom() *ChatRoom {
return &ChatRoom{
clients: make(map[net.Conn]string),
}
}
func (cr *ChatRoom) broadcast(sender net.Conn, msg string) {
cr.mutex.RLock()
defer cr.mutex.RUnlock()
for client, name := range cr.clients {
if sender == client {
continue
}
client.Write([]byte(name + ": " + msg))
}
}
func (cr *ChatRoom) handleConnection(conn net.Conn) {
defer conn.Close()
conn.Write([]byte("请输入你的昵称:"))
name := make([]byte, 1024)
n, _ := conn.Read(name)
cr.mutex.Lock()
cr.clients[conn] = string(name[:n-1])
cr.mutex.Unlock()
username := string(name[:n-1])
cr.broadcast(conn, fmt.Sprintf("系统: %s 加入了聊天室\n", username))
for {
readBuff := make([]byte, 1024)
n, err := conn.Read(readBuff)
if err != nil {
break
}
message := string(readBuff[:n])
if message == "quit\n" {
break
}
cr.broadcast(conn, fmt.Sprintf("%s: %s\n", username, message))
}
cr.mutex.Lock()
delete(cr.clients, conn)
cr.mutex.Unlock()
cr.broadcast(conn, fmt.Sprintf("系统: %s 离开了聊天室", username))
}
func main() {
chatRoom := NewChatRoom()
listener, err := net.Listen("tcp", ":8080")
if err != nil {
panic(err)
}
defer listener.Close()
fmt.Println("启动监听成功")
for {
conn, _ := listener.Accept()
go chatRoom.handleConnection(conn)
}
}goHTTP编程(net/http包、创建HTTP服务器、处理请求与响应)#
什么是HTTP?#
HTTP(HyperText Transfer Protocol,超文本传输协议)是互联网上应用最为广泛的一种网络协议,用于客户端和服务器之间的通信。
核心特点:
- 无状态协议:每次请求都是独立的,服务器不保留之前请求的信息
- 基于请求-响应模型:客户端发起请求,服务器返回响应
- 应用层协议:建立在TCP/IP协议之上
- 明文传输:数据不加密(HTTPS是加密版本)
HTTP请求-响应流程#
客户端 (浏览器/APP) 服务器 (Web服务)
| |
| --- HTTP请求 -------> |
| |
| <--- HTTP响应 -------- |
| |plaintextHTTP请求组成#
GET /api/students HTTP/1.1 ← 请求行(方法 + URL + 协议版本)
Host: localhost:8080 ← 请求头(元数据)
Content-Type: application/json
Authorization: Bearer token123
← 空行分隔
{"name": "张三", "age": 20} ← 请求体(可选,POST/PUT时有)plaintextHTTP方法语义:
GET:获取资源(查)POST:创建资源(增)PUT:更新资源(改)DELETE:删除资源(删)PATCH:部分更新
HTTP响应组成#
HTTP/1.1 200 OK ← 状态行(协议版本 + 状态码 + 描述)
Content-Type: application/json ← 响应头(元数据)
Content-Length: 45
Date: Mon, 23 Oct 2023 08:00:00 GMT
← 空行分隔
{"id": 1, "name": "张三", "age": 20} ← 响应体(数据内容)plaintext常见状态码:
200 OK:请求成功201 Created:资源创建成功400 Bad Request:客户端请求错误401 Unauthorized:未认证404 Not Found:资源不存在500 Internal Server Error:服务器内部错误
net/http包架构#
客户端请求
↓
http.ServeMux (路由 multiplexer)
↓
http.Handler (处理器接口)
↓
具体处理逻辑 (HandlerFunc或自定义Handler)
↓
http.ResponseWriter (写回响应)plaintext核心组件说明#
1. 路由 (ServeMux)
- 负责将不同的URL路径映射到对应的处理器
- 可以创建多个路由,但通常一个应用使用一个主路由
2. 处理器 (Handler)
// Handler是一个接口,只需要实现一个方法
type Handler interface {
ServeHTTP(ResponseWriter, *Request)
}go3. 处理器函数 (HandlerFunc)
// 函数签名,与ServeHTTP方法相同
type HandlerFunc func(ResponseWriter, *Request)go4. 请求对象 (Request)
- 包含客户端的全部请求信息
- 方法、URL、请求头、请求体等
5. 响应写入器 (ResponseWriter)
- 用于构建和发送响应回客户端
- 可以设置状态码、响应头、写入响应体
package main
import (
"fmt"
"net/http"
)
func main() {
mux := http.NewServeMux()
mux.HandleFunc("/students", func(w http.ResponseWriter, r *http.Request) {
w.Header().Add("Content-Type", "application/json")
if r.Method == http.MethodGet {
fmt.Fprintf(w, `[{"id" : 1, "name": "张三", "age": 20}]`)
return
}
http.Error(w, "Method Not Allowed", http.StatusMethodNotAllowed)
})
http.ListenAndServe(":8080", mux)
}
goWebsocket#
WebSocket是什么?
想象一下你和朋友打电话的场景:
- HTTP:像发短信,每次都要建立连接→发送→断开,不能实时对话
- WebSocket:像打电话,一次接通后双方可以随时说话,实现真正的实时对话
WebSocket核心特点:
- 一次握手,长久连接:客户端发起WebSocket请求,服务端同意后建立持久连接
- 双向通信:服务端可以主动向客户端推送数据,不再需要客户端轮询
- 低延迟:避免了HTTP每次请求的头部开销和连接建立时间
- 协议升级:基于HTTP协议升级而来(HTTP 101状态码)
WebSocket握手过程:
客户端请求:
GET /chat HTTP/1.1
Host: server.example.com
Upgrade: websocket
Connection: Upgrade
服务端响应:
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgradeplaintext为什么需要gorilla/websocket库? Go标准库没有提供完整的WebSocket实现,gorilla/websocket是业界最成熟的选择:
- 处理了WebSocket协议细节
- 提供了连接管理、消息读写等高级功能
- 有良好的错误处理和连接状态管理
package main
import (
"fmt"
"log"
"net/http"
"github.com/gorilla/websocket"
)
// 创建WebSocket升级器
var upgrader = websocket.Upgrader{
// 允许所有跨域请求(生产环境应该严格限制)
CheckOrigin: func(r *http.Request) bool {
return true
},
}
func main() {
// 设置WebSocket路由
http.HandleFunc("/ws", handleWebSocket)
// 启动静态文件服务(用于提供HTML页面)
http.Handle("/", http.FileServer(http.Dir("./public")))
fmt.Println("WebSocket服务器启动在 :8080")
fmt.Println("访问 http://localhost:8080 测试聊天功能")
log.Fatal(http.ListenAndServe(":8080", nil))
}
func handleWebSocket(w http.ResponseWriter, r *http.Request) {
// 1. 升级HTTP连接到WebSocket连接
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
log.Println("升级WebSocket失败:", err)
return
}
defer conn.Close() // 确保连接最终会关闭
fmt.Println("新的WebSocket连接建立!")
// 2. 持续监听和处理消息
for {
// 读取客户端发送的消息
messageType, message, err := conn.ReadMessage()
if err != nil {
log.Println("读取消息失败:", err)
break
}
fmt.Printf("收到消息: %s\n", message)
// 3. 向客户端回送消息
response := fmt.Sprintf("服务器回复: 收到你的消息 '%s'", message)
err = conn.WriteMessage(messageType, []byte(response))
if err != nil {
log.Println("发送消息失败:", err)
break
}
}
fmt.Println("WebSocket连接关闭")
}go