nsq源碼(9) nsqlookupd與nsqd交互

nsqlookupd與nsq交互

nsqd帶參數(shù)啟動

  • 除了接收pub發(fā)布的topic,還可以通過硬盤備份的文件恢復(fù)創(chuàng)建topic
  • 在nsqd啟動時除了開啟tcp,http與queueScanLoop超時檢測線程批狱,還有開啟一個lookupLoop線程去注冊nsqlookupd
func (p *program) Start() error {
    opts := nsqd.NewOptions()

    // 讀取參數(shù)設(shè)置Options對象
    flagSet := nsqdFlagSet(opts)
    flagSet.Parse(os.Args[1:])

    options.Resolve(opts, flagSet, cfg)
    // 使用Options對象創(chuàng)建nsqd對象
    nsqd := nsqd.New(opts)

    nsqd.Main()
    ...
}

func (n *NSQD) Main() {
    // 恢復(fù)創(chuàng)建備份文件中的topic
    err := nsqd.LoadMetadata()
    if err != nil {
        log.Fatalf("ERROR: %s", err.Error())
    }

    // 超時消息檢索和處理任務(wù)
    n.waitGroup.Wrap(n.queueScanLoop)

    // 根據(jù)參數(shù)選擇注冊中心nsqlookupd
    n.waitGroup.Wrap(n.lookupLoop)
    if n.getOpts().StatsdAddress != "" {
        n.waitGroup.Wrap(n.statsdLoop)
    }
}

NSQD.LoadMetadata() 備份文件創(chuàng)建topic

  • 在創(chuàng)建topic時除了開啟messagePump線程接收memoryMsgChan隊(duì)列
  • 還會nsqd.Notify(t)通過topic.notifyChan隊(duì)列通知lookupLoop線程去注冊該topic
func (n *NSQD) LoadMetadata() error {
    // 讀取nsqd.dat硬盤備份文件
    fn := newMetadataFile(n.getOpts())

    // 根據(jù)備份文件創(chuàng)建topic
    for _, t := range m.Topics {
        if !protocol.IsValidTopicName(t.Name) {
            n.logf(LOG_WARN, "skipping creation of invalid topic %s", t.Name)
            continue
        }
        topic := n.GetTopic(t.Name)
        if t.Paused {
            topic.Pause()
        }
        for _, c := range t.Channels {
            if !protocol.IsValidChannelName(c.Name) {
                n.logf(LOG_WARN, "skipping creation of invalid channel %s", c.Name)
                continue
            }
            channel := topic.GetChannel(c.Name)
            if c.Paused {
                channel.Pause()
            }
        }
        topic.Start()
    }
    return nil
}

func NewTopic(topicName string, ctx *context, deleteCallback func(*Topic)) *Topic {
    t.waitGroup.Wrap(t.messagePump)

    // 退出或者通知nsqdlookup進(jìn)行注冊操作
    // 通過topic.notifyChan隊(duì)列通知lookupLoop線程去注冊該topic
    t.ctx.nsqd.Notify(t)

    return t
}

NSQD.lookupLoop 接收通知進(jìn)行topic/channel注冊

  • 在nsqd啟動時開啟的lookupLoop線程會循環(huán)處理notifyChan隊(duì)列發(fā)來的消息
// 為當(dāng)前nsqd綁定nsqlookupd
func (n *NSQD) lookupLoop() {
    for {
        if connect {
            for _, host := range n.getOpts().NSQLookupdTCPAddresses {
                if in(host, lookupAddrs) {
                    continue
                }
                n.logf(LOG_INFO, "LOOKUP(%s): adding peer", host)
                // 讀取參數(shù)
                // LOOKUP connecting to 127.0.0.1:4160
                // LOOKUPD(127.0.0.1:4160): peer info {TCPPort:4160 HTTPPort:4161 Version:1.1.0 BroadcastAddress:sz-linrundong}
                lookupPeer := newLookupPeer(host, n.getOpts().MaxBodySize, n.logf,
                    connectCallback(n, hostname))
                // 進(jìn)行連接nsqlooupd操作
                lookupPeer.Command(nil) // start the connection
                lookupPeers = append(lookupPeers, lookupPeer)
                lookupAddrs = append(lookupAddrs, host)
            }
            n.lookupPeers.Store(lookupPeers)
            connect = false
        }

        select {
        // 通過NSQD.notifyChan隊(duì)列獲取nsqlookupd發(fā)來的消息
        case val := <-n.notifyChan:
            var cmd *nsq.Command
            var branch string

            // 斷言解析隊(duì)列消息
            switch val.(type) {
            case *Channel:
                // notify all nsqlookupds that a new channel exists, or that it's removed
                branch = "channel"
                channel := val.(*Channel)
                if channel.Exiting() == true {
                    cmd = nsq.UnRegister(channel.topicName, channel.name)
                } else {
                    //拼裝注冊命令
                    cmd = nsq.Register(channel.topicName, channel.name)
                }
            case *Topic:
                // notify all nsqlookupds that a new topic exists, or that it's removed
                branch = "topic"
                topic := val.(*Topic)
                if topic.Exiting() == true {
                    cmd = nsq.UnRegister(topic.name, "")
                } else {
                    cmd = nsq.Register(topic.name, "")
                }
            }

            // 向每個lookupd發(fā)送請求命令cmd
            for _, lookupPeer := range lookupPeers {
                n.logf(LOG_INFO, "LOOKUPD(%s): %s %s", lookupPeer, branch, cmd)
                _, err := lookupPeer.Command(cmd)
                if err != nil {
                    n.logf(LOG_ERROR, "LOOKUPD(%s): %s - %s", lookupPeer, cmd, err)
                }
            }
        case <-n.exitChan:
            goto exit
        }
    }

exit:
    n.logf(LOG_INFO, "LOOKUP: closing")
}

向nsqlookupd發(fā)送注冊請求

  • 發(fā)送選項(xiàng)以及序列化channel训挡,獲取響應(yīng)
func (lp *lookupPeer) Command(cmd *nsq.Command) ([]byte, error) {
    initialState := lp.state
    if lp.state != stateConnected {
        err := lp.Connect()
        if err != nil {
            return nil, err
        }
        lp.state = stateConnected
        _, err = lp.Write(nsq.MagicV1)
        if err != nil {
            lp.Close()
            return nil, err
        }
        if initialState == stateDisconnected {
            lp.connectCallback(lp)
        }
        if lp.state != stateConnected {
            return nil, fmt.Errorf("lookupPeer connectCallback() failed")
        }
    }
    if cmd == nil {
        return nil, nil
    }
    _, err := cmd.WriteTo(lp)
    if err != nil {
        lp.Close()
        return nil, err
    }
    resp, err := readResponseBounded(lp, lp.maxBodySize)
    if err != nil {
        lp.Close()
        return nil, err
    }
    return resp, nil
}
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
  • 序言:七十年代末娩嚼,一起剝皮案震驚了整個濱河市,隨后出現(xiàn)的幾起案子,更是在濱河造成了極大的恐慌谬以,老刑警劉巖,帶你破解...
    沈念sama閱讀 211,948評論 6 492
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件圆雁,死亡現(xiàn)場離奇詭異忍级,居然都是意外死亡,警方通過查閱死者的電腦和手機(jī)摸柄,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,371評論 3 385
  • 文/潘曉璐 我一進(jìn)店門颤练,熙熙樓的掌柜王于貴愁眉苦臉地迎上來,“玉大人驱负,你說我怎么就攤上這事嗦玖。” “怎么了跃脊?”我有些...
    開封第一講書人閱讀 157,490評論 0 348
  • 文/不壞的土叔 我叫張陵宇挫,是天一觀的道長。 經(jīng)常有香客問我酪术,道長器瘪,這世上最難降的妖魔是什么? 我笑而不...
    開封第一講書人閱讀 56,521評論 1 284
  • 正文 為了忘掉前任绘雁,我火速辦了婚禮橡疼,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘庐舟。我一直安慰自己欣除,他們只是感情好,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,627評論 6 386
  • 文/花漫 我一把揭開白布挪略。 她就那樣靜靜地躺著历帚,像睡著了一般。 火紅的嫁衣襯著肌膚如雪杠娱。 梳的紋絲不亂的頭發(fā)上挽牢,一...
    開封第一講書人閱讀 49,842評論 1 290
  • 那天,我揣著相機(jī)與錄音摊求,去河邊找鬼禽拔。 笑死,一個胖子當(dāng)著我的面吹牛睹簇,可吹牛的內(nèi)容都是我干的奏赘。 我是一名探鬼主播,決...
    沈念sama閱讀 38,997評論 3 408
  • 文/蒼蘭香墨 我猛地睜開眼太惠,長吁一口氣:“原來是場噩夢啊……” “哼磨淌!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起凿渊,我...
    開封第一講書人閱讀 37,741評論 0 268
  • 序言:老撾萬榮一對情侶失蹤梁只,失蹤者是張志新(化名)和其女友劉穎缚柳,沒想到半個月后,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體搪锣,經(jīng)...
    沈念sama閱讀 44,203評論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡秋忙,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,534評論 2 327
  • 正文 我和宋清朗相戀三年,在試婚紗的時候發(fā)現(xiàn)自己被綠了构舟。 大學(xué)時的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片灰追。...
    茶點(diǎn)故事閱讀 38,673評論 1 341
  • 序言:一個原本活蹦亂跳的男人離奇死亡,死狀恐怖狗超,靈堂內(nèi)的尸體忽然破棺而出弹澎,到底是詐尸還是另有隱情,我是刑警寧澤努咐,帶...
    沈念sama閱讀 34,339評論 4 330
  • 正文 年R本政府宣布苦蒿,位于F島的核電站,受9級特大地震影響渗稍,放射性物質(zhì)發(fā)生泄漏佩迟。R本人自食惡果不足惜,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 39,955評論 3 313
  • 文/蒙蒙 一竿屹、第九天 我趴在偏房一處隱蔽的房頂上張望报强。 院中可真熱鬧,春花似錦拱燃、人聲如沸躺涝。這莊子的主人今日做“春日...
    開封第一講書人閱讀 30,770評論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽。三九已至夯膀,卻和暖如春诗充,著一層夾襖步出監(jiān)牢的瞬間,已是汗流浹背诱建。 一陣腳步聲響...
    開封第一講書人閱讀 32,000評論 1 266
  • 我被黑心中介騙來泰國打工蝴蜓, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留,地道東北人俺猿。 一個月前我還...
    沈念sama閱讀 46,394評論 2 360
  • 正文 我出身青樓茎匠,卻偏偏與公主長得像,于是被迫代替她去往敵國和親押袍。 傳聞我的和親對象是個殘疾皇子诵冒,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,562評論 2 349

推薦閱讀更多精彩內(nèi)容