nsq源码(2) nsqd TcpServer封装

启动流程

  • 启动方式与nsqlookupd相同使用svg
  • 创建两个协程分别监听HTTP和TCP端口
// nsqd.main()
func main() {
    prg := &program{}
    if err := svc.Run(prg, syscall.SIGINT, syscall.SIGTERM); err != nil {
        log.Fatal(err)
    }
    
    // 初始化
    if err = prg.Init(service); err != nil {
        return err
    }
    
    // 启动业务进程
    err := prg.Start()
    if err != nil {
        return err
    }
    
    // 阻塞等待退出信号
    signalChan := make(chan os.Signal, 1)
    svg.signalNotify(signalChan, ws.signals...)
    <-signalChan
    
    // 退出
    err = prg.Stop()
}

Main线程与nsqlookupd差异

  • nsqlookupd的HTTP线程与消费者通讯, TCP端口作为Server与nsq 通讯(长连接)
  • nsqd的TCP端口接收客户端(生产者)消息,HTTP端口提供HTTP API接收客户端(生产者)消息
  • nsqd与nsqlookupd交互方式以及nsqd注册在后续章节
func (n *NSQD) Main() {
    tcpServer := &tcpServer{ctx: ctx}
    n.waitGroup.Wrap(func() {
        // 协程连接客户端(生产者)
        protocol.TCPServer(n.tcpListener, tcpServer, n.logf)
    })
    
    // 监听端口提供HTTP API
    httpServer := newHTTPServer(ctx, false, n.getOpts().TLSRequired == TLSRequired)
    n.waitGroup.Wrap(func() {
        http_api.Serve(n.httpListener, httpServer, "HTTP", n.logf)
    })
    if n.tlsConfig != nil && n.getOpts().HTTPSAddress != "" {
        httpsServer := newHTTPServer(ctx, true, true)
        n.waitGroup.Wrap(func() {
            http_api.Serve(n.httpsListener, httpsServer, "HTTPS", n.logf)
        })
    }
    
    // 超时消息检索和处理任务
    n.waitGroup.Wrap(n.queueScanLoop)
    n.waitGroup.Wrap(n.queueScanLoop)
    n.waitGroup.Wrap(n.lookupLoop)
    if n.getOpts().StatsdAddress != "" {
        n.waitGroup.Wrap(n.statsdLoop)
    }
}

可复用的Tcp Server模式

  • TcpServer封装
  • 循环阻塞accept()
  • 当退出信号触发svg设置的channel,将关闭listener,从而将在Handle处理完毕之后才会退出
  • 每来一个client都开启一个Handler协程
func TCPServer(listener net.Listener, handler TCPHandler, logf lg.AppLogFunc) {
    logf(lg.INFO, "TCP: listening on %s", listener.Addr())

    for {
        clientConn, err := listener.Accept()
        if err != nil {
            if nerr, ok := err.(net.Error); ok && nerr.Temporary() {
                logf(lg.WARN, "temporary Accept() failure - %s", err)
                runtime.Gosched()
                continue
            }
            // theres no direct way to detect this error because it is not exposed
            if !strings.Contains(err.Error(), "use of closed network connection") {
                logf(lg.ERROR, "listener.Accept() - %s", err)
            }
            break
        }
        go handler.Handle(clientConn)
    }

    logf(lg.INFO, "TCP: closing %s", listener.Addr())
}
数据读取,使用protocol拓展的模式
func (p *tcpServer) Handle(clientConn net.Conn) {
    p.ctx.nsqd.logf(LOG_INFO, "TCP: new client(%s)", clientConn.RemoteAddr())

    // 读取4字节的协议版本
    buf := make([]byte, 4)
    _, err := io.ReadFull(clientConn, buf)
    ...
    protocolMagic := string(buf)
    ...

    var prot protocol.Protocol
    switch protocolMagic {
    case "  V2":
        // 用报文与nsqd本身信息初始化一个protocolV2
        prot = &protocolV2{ctx: p.ctx}
    default:
        clientConn.Close()
        ...
        return
    }

    // protocolV2.IOLoop()来处理客户端的业务内容
    err = prot.IOLoop(clientConn)
    ...
}

nsqd 对Client的处理

通过对应client的Handle协程的IOLoop()、
func (p *protocolV2) IOLoop(conn net.Conn) error {
    // 创建新client处理类
    clientID := atomic.AddInt64(&p.ctx.nsqd.clientIDSequence, 1)
    client := newClientV2(clientID, conn, p.ctx)
    p.ctx.nsqd.AddClient(client.ID, client)

    // messagePump协程完成初始化向chan发送消息
    // TODO messagePump()作用未知
    messagePumpStartedChan := make(chan bool)
    go p.messagePump(client, messagePumpStartedChan)
    <-messagePumpStartedChan

    for {
        // 心跳如果有设置,则设置读超时为两倍时间
        if client.HeartbeatInterval > 0 {
            client.SetReadDeadline(time.Now().Add(client.HeartbeatInterval * 2))
        } else {
            client.SetReadDeadline(zeroTime)
        }

        // 分割命令
        line, err = client.Reader.ReadSlice('\n')
        // 处理win换行
        line = line[:len(line)-1]
        if len(line) > 0 && line[len(line)-1] == '\r' {
            line = line[:len(line)-1]
        }
        params := bytes.Split(line, separatorBytes)
        p.ctx.nsqd.logf(LOG_DEBUG, "PROTOCOL(V2): [%s] %s", client, params)

        var response []byte
        // Exec()处理实际参数
        response, err = p.Exec(client, params)

        if response != nil {
            // 发送响应
            err = p.Send(client, frameTypeResponse, response)
        }
    }

    p.ctx.nsqd.logf(LOG_INFO, "PROTOCOL(V2): [%s] exiting ioloop", client)
    conn.Close()
    close(client.ExitChan)

    // 移除Client
    p.ctx.nsqd.RemoveClient(client.ID)
    return err
}

Send()发送响应

func (p *protocolV2) Send(client *clientV2, frameType int32, data []byte) error {
    client.writeLock.Lock()

    var zeroTime time.Time
    // 如果设置了心跳则设置写超时
    if client.HeartbeatInterval > 0 {
        client.SetWriteDeadline(time.Now().Add(client.HeartbeatInterval))
    } else {
        client.SetWriteDeadline(zeroTime)
    }

    // protocol.SendFramedResponse()
    SendFramedResponse := func(client.Writer, frameType, data){
        beBuf := make([]byte, 4)
        size := uint32(len(data)) + 4
        // 转换写入数据大小
        binary.BigEndian.PutUint32(beBuf, size)
        n, err := w.Write(beBuf)
        if err != nil {
            return n, err
        }
        // 转换写入数据
        binary.BigEndian.PutUint32(beBuf, uint32(frameType))
        n, err = w.Write(beBuf)
        if err != nil {
            return n + 4, err
        }
    
        n, err = w.Write(data)
        return n + 8, err
    }

    if frameType != frameTypeMessage {
        err = client.Flush()
    }

    client.writeLock.Unlock()

    return err
}
最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 219,635评论 6 508
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 93,628评论 3 396
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 165,971评论 0 356
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,986评论 1 295
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 68,006评论 6 394
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,784评论 1 307
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,475评论 3 420
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 39,364评论 0 276
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,860评论 1 317
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 38,008评论 3 338
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 40,152评论 1 351
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,829评论 5 346
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,490评论 3 331
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 32,035评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 33,156评论 1 272
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 48,428评论 3 373
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 45,127评论 2 356

推荐阅读更多精彩内容

  • 参考教程 为什么用Tornado? 异步编程原理 服务器同时要对许多客户端提供服务,他的性能至关重要。而服务器端的...
    内心强大的Jim阅读 6,215评论 1 8
  • 目录 一、开启线程的两种方式 在python中开启线程要导入threading,它与开启进程所需要导入的模块mul...
    CaiGuangyin阅读 2,408评论 1 16
  • 计算机网络概述 网络编程的实质就是两个(或多个)设备(例如计算机)之间的数据传输。 按照计算机网络的定义,通过一定...
    蛋炒饭_By阅读 1,227评论 0 10
  • 1. 开始 开篇罗嗦了一大堆,终于开始进入正题了。golang的优秀源码很多,比如杀手级应用docker,goog...
    kekemuyu阅读 2,823评论 1 4
  • 网络 理论模型,分为七层物理层数据链路层传输层会话层表示层应用层 实际应用,分为四层链路层网络层传输层应用层 IP...
    FlyingLittlePG阅读 778评论 0 0