udpHandler
# 1.入口
# 1.1.entrypoint
traefik执行setupServer()初始化服务调用server.NewUDPEntryPoints()初始化udp入口,udp流量会根据入口分发给handler处理。// NewUDPEntryPoints returns all the UDP entry points, keyed by name. func NewUDPEntryPoints(cfg static.EntryPoints) (UDPEntryPoints, error) { ... // 遍历入口配置 for entryPointName, entryPoint := range cfg { ... if entryPoint.GetProtocol() != "udp" { continue } // 构造入口 ep, err := NewUDPEntryPoint(entryPoint) ... entryPoints[entryPointName] = ep } return entryPoints, nil } // NewUDPEntryPoint returns a UDP entry point. func NewUDPEntryPoint(cfg *static.EntryPoint) (*UDPEntryPoint, error) { listenConfig := newListenConfig(cfg) // 初始化udp listener listener, err := udp.Listen(listenConfig, "udp", cfg.GetAddress(), time.Duration(cfg.UDP.Timeout)) ... return &UDPEntryPoint{listener: listener, switcher: &udp.HandlerSwitcher{}, transportConfiguration: cfg.Transport}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28注意
udp switcher是服务核心,入口流量会经过switcher分发给handler处理
# 1.2.listener
udp.Listen实现了基于udp的伪连接监听器,使udp看起来像tcp一样可以被accept接受多个连接,其实是listen内部抽象实现。// Listen creates a new listener. func Listen(listenConfig net.ListenConfig, network, address string, timeout time.Duration) (*Listener, error) { ... // 打开udp socket packetConn, err := listenConfig.ListenPacket(context.Background(), network, address) ... // 类型断言为udp连接 pConn, ok := packetConn.(*net.UDPConn) ... // 构造listener实例 l := &Listener{ pConn: pConn, acceptCh: make(chan *Conn), conns: make(map[string]*Conn), accepting: true, timeout: timeout, } // 异步读取udp socket数据 // 将无连接数据流拆分为逻辑连接 go l.readLoop() return l, nil } // 根据连接划分读取到的数据 func (l *Listener) readLoop() { for { // 每次读64K数据 buf := make([]byte, maxDatagramSize) // 向底层udp socket读取数据包 n, raddr, err := l.pConn.ReadFrom(buf) ... // 获取远程地址连接 conn, err := l.getConn(raddr) ... select { // 数据推送 case conn.receiveCh <- buf[:n]: // 连接关闭,丢弃本次数据 case <-conn.doneCh: continue } } } // getConn returns the ongoing session with raddr if it exists, or creates a new one otherwise. func (l *Listener) getConn(raddr net.Addr) (*Conn, error) { l.mu.Lock() defer l.mu.Unlock() // 先获取缓存的 conn, ok := l.conns[raddr.String()] if ok { return conn, nil } ... // 新建连接 conn = l.newConn(raddr) // 缓存连接 l.conns[raddr.String()] = conn // 连接推送给accept l.acceptCh <- conn // 读请求处理 go conn.readLoop() return conn, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74注意
listener的伪监听实现本质上是读取udp socket数据及暂存,将抽象连接推给listener.Accept(),之后的读操作就是从缓存中获取数据
# 1.3.start
udp entrypoint启动会基于listener不停监听伪连接,执行switcher.ServeUDP()将连接分发给handler以获取请求数据及处理。// Start commences the listening for all the entry points. func (eps UDPEntryPoints) Start() { // 根据入口依次启动 for entryPointName, ep := range eps { go ep.Start(ctx) } } // Start commences the listening for ep. func (ep *UDPEntryPoint) Start(ctx context.Context) { for { // 获取连接 conn, err := ep.listener.Accept() ... // 请求处理 go ep.switcher.ServeUDP(conn) } } // ServeUDP implements the Handler interface. func (s *HandlerSwitcher) ServeUDP(conn *Conn) { // 获取udp handler handler := s.handler.Get() h, ok := handler.(Handler) if ok { // 处理流量 h.ServeUDP(conn) } else { conn.Close() } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33注意
udp switcher较为特殊,入口仅关联一个handler,所以没有router对象匹配路由,直接由switcher分发请求到service handler
# 2.路由
# 2.1.switchRouter
watcher监听配置变更会触发switchRouter构造及切换路由,udp相关的路由构造及切换就是routerFactory和udp switcher完成的。func switchRouter(...) func(conf dynamic.Configuration) { return func(conf dynamic.Configuration) { ... // 构造 routers, udpRouters := routerFactory.CreateRouters(rtConf) ... // 切换 serverEntryPointsUDP.Switch(udpRouters) } } // CreateRouters creates new TCPRouters and UDPRouters. func (f *RouterFactory) CreateRouters(rtConf *runtime.Configuration) (map[string]*tcprouter.Router, map[string]udp.Handler) { ... // UDP rtUDPManager := udprouter.NewManager(rtConf, svcUDPManager) // 构造handler routersUDP := rtUDPManager.BuildHandlers(ctx, f.entryPointsUDP) ... return routersTCP, routersUDP } // Switch swaps out all the given handlers in their associated entrypoints. func (eps UDPEntryPoints) Switch(handlers map[string]udp.Handler) { for epName, handler := range handlers { // 基于switcher切换handler if ep, ok := eps[epName]; ok { ep.Switch(handler) continue } } } // Switch replaces ep's handler with the one given as argument. func (ep *UDPEntryPoint) Switch(handler udp.Handler) { ep.switcher.Switch(handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41注意
udp switcher切换handler本质上是换了字段值,旧的流量不受影响,核心还是handler构造
# 2.2.buildHandlers
buildHandlers()用于构造udp handler,每个入口仅关联一个udp handler,不涉及router及AST语法树匹配路由,相对tcp简单一些。// BuildHandlers builds the handlers for the given entrypoints. func (m *Manager) BuildHandlers(rootCtx context.Context, entryPoints []string) map[string]udp.Handler { // 获取路由定义 entryPointsRouters := m.getUDPRouters(rootCtx, entryPoints) ... // 遍历入口 for _, entryPointName := range entryPoints { // 获取入口关联路由配置 routers := entryPointsRouters[entryPointName] ... // 构造udp handler handlers := m.buildEntryPointHandlers(ctx, routers) if len(handlers) > 0 { // As UDP support only one router per entrypoint, we only take the first one. entryPointHandlers[entryPointName] = handlers[0] } } return entryPointHandlers } func (m *Manager) buildEntryPointHandlers(ctx context.Context, configs map[string]*runtime.UDPRouterInfo) []udp.Handler { ... // router名称根据字典倒序 sort.Slice(rtNames, func(i, j int) bool { return rtNames[i] > rtNames[j] }) ... // 遍历router for _, routerName := range rtNames { routerConfig := configs[routerName] ... // 构造udp handler handler, err := m.serviceManager.BuildUDP(ctxRouter, routerConfig.Service) ... handlers = append(handlers, handler) } return handlers } // udp service handler func (m *Manager) BuildUDP(rootCtx context.Context, serviceName string) (udp.Handler, error) { ... // 获取service配置 conf, ok := m.configs[serviceQualifiedName] ... switch { // LB类型 case conf.LoadBalancer != nil: loadBalancer := udp.NewWRRLoadBalancer() // 打乱遍历后端池 for index, server := range shuffle(conf.LoadBalancer.Servers, m.rand) { ... // 创建udp后端代理 handler, err := udp.NewProxy(server.Address) ... // 注册到LB loadBalancer.AddServer(handler) } return loadBalancer, nil // 加权负载均衡 case conf.Weighted != nil: loadBalancer := udp.NewWRRLoadBalancer() // 打乱遍历service for _, service := range shuffle(conf.Weighted.Services, m.rand) { // 递归解析service LB handler, err := m.BuildUDP(ctx, service.Name) ... // 注册到WRR负载均衡器 loadBalancer.AddWeightedServer(handler, service.Weight) } return loadBalancer, nil default: ... return nil, err } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90注意
udp handler本质上只有loadbalancer类型,后端以udp proxy形式存在
# 3.流量处理
# 3.1.balancer
balancer是udp handler的核心实现,listener.Accept()收到连接会执行switcher.ServeUDP(),switcher将连接转给balancer。// ServeUDP forwards the connection to the right service. func (b *WRRLoadBalancer) ServeUDP(conn *Conn) { b.lock.Lock() // 选一个后端 next, err := b.next() b.lock.Unlock() ... // 执行proxy.ServeUDP() next.ServeUDP(conn) } func (b *WRRLoadBalancer) next() (Handler, error) { ... // 获取服务最大权重 max := b.maxWeight() ... // 所有服务权重最大公约数 gcd := b.weightGcd() for { // 轮询下一个服务器(index初始为-1) b.index = (b.index + 1) % len(b.servers) // 执行一轮 if b.index == 0 { // 阈值减GCD b.currentWeight -= gcd // 阈值每轮重置到最大 if b.currentWeight <= 0 { b.currentWeight = max } } // 获取该轮后端 srv := b.servers[b.index] // 满足阈值则返回 if srv.weight >= b.currentWeight { return srv, nil } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42注意
balancer的实现与tcp类似,根据权重算法选出后端后执行next.ServeUDP(),这里的后端可能是balancer,也可能是udp proxy
# 3.2.proxy
udp proxy会创建面向后端的udp连接,通过零拷贝将conn流和backend流绑定,以转发请求数据到后端服务及响应结果到客户端。// ServeUDP implements the Handler interface. func (p *Proxy) ServeUDP(conn *Conn) { // needed because of e.g. server.trackedConnection defer conn.Close() // 创建后端连接 connBackend, err := net.Dial("udp", p.target) ... // maybe not needed, but just in case defer connBackend.Close() ... // 流拷贝 go connCopy(conn, connBackend, errChan) go connCopy(connBackend, conn, errChan) ... <-errChan } // conn Read重写 func (c *Conn) Read(p []byte) (int, error) { select { // 读请求 case c.readCh <- p: n := <-c.sizeCh c.muActivity.Lock() c.lastActivity = time.Now() c.muActivity.Unlock() return n, nil case <-c.doneCh: return 0, io.EOF } } // conn Write重写 func (c *Conn) Write(p []byte) (n int, err error) { c.muActivity.Lock() c.lastActivity = time.Now() c.muActivity.Unlock() // 基于udp socket将结果写给客户端 return c.listener.pConn.WriteTo(p, c.rAddr) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46注意
proxy.ServeUDP()处理类似tcp代理,基于零拷贝复制conn与backend的写入写出流,请求数据直接转给后端,响应结果直接写给客户端