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

如何搭建无端口暴露的分布式转发代理用于网页爬取?

技术方案实现指引

核心架构确认

你的思路完全可行:客户端主动与中心化服务器建立WebSocket长连接(无需端口暴露),服务器作为常规HTTP代理接收请求方的请求,将请求转发给随机客户端执行,最终把响应原路返回。这是典型的反向代理+WebSocket隧道模式,完美解决无端口暴露的问题。

现有代码的核心问题

你尝试用Hijack直接转发TCP流的思路不适合WebSocket场景:

  1. WebSocket是帧化协议,无法直接用io.Copy转发原始HTTP流
  2. 缺少客户端连接池管理,无法实现请求的随机分配
  3. 没有处理请求与响应的关联逻辑,服务器无法将客户端返回的响应匹配到对应的请求方

具体实现代码示例

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

关键实现要点

  1. 连接池管理:服务器维护在线客户端列表,随机分配请求实现负载均衡
  2. 请求序列化:用httputil.DumpRequest/Response将HTTP请求/响应转成字节数组,通过WebSocket的JSON消息传递
  3. 请求关联:用唯一ID绑定请求与响应,确保服务器能正确匹配返回结果
  4. 超时与错误处理:全链路添加超时控制和错误捕获,避免阻塞或崩溃
  5. 权限校验:生产环境需为WebSocket连接添加身份验证,防止非法客户端接入

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 04:55:05