如何搭建无端口暴露的分布式转发代理用于网页爬取?
技术方案实现指引
核心架构确认
你的思路完全可行:客户端主动与中心化服务器建立WebSocket长连接(无需端口暴露),服务器作为常规HTTP代理接收请求方的请求,将请求转发给随机客户端执行,最终把响应原路返回。这是典型的反向代理+WebSocket隧道模式,完美解决无端口暴露的问题。
现有代码的核心问题
你尝试用Hijack直接转发TCP流的思路不适合WebSocket场景:
- WebSocket是帧化协议,无法直接用
io.Copy转发原始HTTP流 - 缺少客户端连接池管理,无法实现请求的随机分配
- 没有处理请求与响应的关联逻辑,服务器无法将客户端返回的响应匹配到对应的请求方
具体实现代码示例
1. 中心化服务器代码
package main import ( "bytes" "encoding/json" "fmt" "math/rand" "net/http" "net/http/httputil" "sync" "time" "github.com/gorilla/websocket" ) var upgrader = websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true // 生产环境需替换为实际权限校验逻辑 }, } // 客户端连接池 var clientPool = struct { sync.RWMutex conns []*websocket.Conn }{} // 代理请求消息结构 type ProxyRequest struct { ID string `json:"id"` RawReq []byte `json:"raw_req"` } // 代理响应消息结构 type ProxyResponse struct { ID string `json:"id"` RawResp []byte `json:"raw_resp"` Error string `json:"error,omitempty"` } // 请求关联映射 var reqMap = sync.Map{} func main() { rand.Seed(time.Now().UnixNano()) // WebSocket端点:接收客户端连接 http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) { conn, err := upgrader.Upgrade(w, r, nil) if err != nil { fmt.Println("WebSocket升级失败:", err) return } defer conn.Close() // 将客户端加入连接池 clientPool.Lock() clientPool.conns = append(clientPool.conns, conn) clientPool.Unlock() fmt.Printf("新客户端上线,当前在线数:%d\n", len(clientPool.conns)) // 持续读取客户端返回的响应 for { var resp ProxyResponse if err := conn.ReadJSON(&resp); err != nil { // 客户端断开,从连接池移除 clientPool.Lock() for i, c := range clientPool.conns { if c == conn { clientPool.conns = append(clientPool.conns[:i], clientPool.conns[i+1:]...) break } } clientPool.Unlock() fmt.Printf("客户端下线,当前在线数:%d\n", len(clientPool.conns)) break } // 将响应转发给对应的请求方 if respChan, ok := reqMap.Load(resp.ID); ok { respChan.(chan<- ProxyResponse) <- resp reqMap.Delete(resp.ID) } } }) // HTTP代理端点:接收请求方的代理请求 http.HandleFunc("/", proxyHandler) fmt.Println("服务器启动,监听端口8080") if err := http.ListenAndServe(":8080", nil); err != nil { panic(err) } } func proxyHandler(w http.ResponseWriter, r *http.Request) { // 检查是否有可用客户端 clientPool.RLock() clientCount := len(clientPool.conns) if clientCount == 0 { clientPool.RUnlock() w.WriteHeader(http.StatusServiceUnavailable) fmt.Fprint(w, "无可用代理客户端") return } // 随机选择一个客户端 targetConn := clientPool.conns[rand.Intn(clientCount)] clientPool.RUnlock() // 序列化原始HTTP请求 rawReq, err := httputil.DumpRequest(r, true) if err != nil { w.WriteHeader(http.StatusBadGateway) fmt.Fprint(w, "序列化请求失败:", err) return } // 生成唯一请求ID reqID := fmt.Sprintf("%d", time.Now().UnixNano()) // 创建响应通道 respChan := make(chan ProxyResponse, 1) reqMap.Store(reqID, respChan) defer reqMap.Delete(reqID) defer close(respChan) // 将请求发送给客户端 if err := targetConn.WriteJSON(ProxyRequest{ID: reqID, RawReq: rawReq}); err != nil { w.WriteHeader(http.StatusBadGateway) fmt.Fprint(w, "转发请求到客户端失败:", err) return } // 等待响应或超时 select { case resp := <-respChan: if resp.Error != "" { w.WriteHeader(http.StatusBadGateway) fmt.Fprint(w, "客户端执行请求出错:", resp.Error) return } // 将响应写回请求方 w.WriteHeader(http.StatusOK) w.Write(resp.RawResp) case <-time.After(30 * time.Second): w.WriteHeader(http.StatusGatewayTimeout) fmt.Fprint(w, "请求超时") } }
2. 客户端代码
package main import ( "bytes" "encoding/json" "fmt" "net/http" "net/http/httputil" "github.com/gorilla/websocket" ) // 与服务器一致的消息结构 type ProxyRequest struct { ID string `json:"id"` RawReq []byte `json:"raw_req"` } type ProxyResponse struct { ID string `json:"id"` RawResp []byte `json:"raw_resp"` Error string `json:"error,omitempty"` } func main() { // 连接中心化服务器的WebSocket端点 conn, _, err := websocket.DefaultDialer.Dial("ws://localhost:8080/ws", nil) if err != nil { panic("连接服务器失败:" + err.Error()) } defer conn.Close() fmt.Println("已连接到中心化服务器") // 持续处理服务器发来的代理请求 for { var req ProxyRequest if err := conn.ReadJSON(&req); err != nil { fmt.Println("读取服务器请求失败:", err) break } // 解析原始HTTP请求 httpReq, err := http.ReadRequest(bytes.NewReader(req.RawReq)) if err != nil { sendError(conn, req.ID, "解析请求失败:"+err.Error()) continue } // 执行请求到目标服务器 client := &http.Client{} resp, err := client.Do(httpReq) if err != nil { sendError(conn, req.ID, "发送请求到目标服务器失败:"+err.Error()) continue } defer resp.Body.Close() // 序列化响应 rawResp, err := httputil.DumpResponse(resp, true) if err != nil { sendError(conn, req.ID, "序列化响应失败:"+err.Error()) continue } // 将响应发回服务器 if err := conn.WriteJSON(ProxyResponse{ID: req.ID, RawResp: rawResp}); err != nil { fmt.Println("发送响应到服务器失败:", err) } } } func sendError(conn *websocket.Conn, reqID, errMsg string) { _ = conn.WriteJSON(ProxyResponse{ID: reqID, Error: errMsg}) }
关键实现要点
- 连接池管理:服务器维护在线客户端列表,随机分配请求实现负载均衡
- 请求序列化:用
httputil.DumpRequest/Response将HTTP请求/响应转成字节数组,通过WebSocket的JSON消息传递 - 请求关联:用唯一ID绑定请求与响应,确保服务器能正确匹配返回结果
- 超时与错误处理:全链路添加超时控制和错误捕获,避免阻塞或崩溃
- 权限校验:生产环境需为WebSocket连接添加身份验证,防止非法客户端接入
内容的提问来源于stack exchange,提问作者Preston Harrison
相关产品推荐
相关产品推荐

