tmq3mnodeSwitch.py 10.9 KB
Newer Older
P
plum-lihui 已提交
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

from ntpath import join
import taos
import sys
import time
import socket
import os
import threading

from util.log import *
from util.sql import *
from util.cases import *
from util.dnodes import *
from util.common import *
from util.cluster import *

sys.path.append("./7-tmq")
from tmqCommon import *

class TDTestCase:
    def __init__(self):
        self.dnodes = 5
        self.mnodes = 3
        self.idIndex = 0
        self.roleIndex = 2
        self.mnodeStatusIndex = 3
        self.mnodeEpIndex = 1
        self.dnodeStatusIndex = 4
        self.mnodeCheckCnt    = 10
        self.host = socket.gethostname()
        self.startPort = 6030
        self.portStep = 100
        self.dnodeOfLeader = 0

35
    def init(self, conn, logSql, replicaVar=1):
P
plum-lihui 已提交
36 37 38
        tdLog.debug(f"start to excute {__file__}")
        tdSql.init(conn.cursor())
        #tdSql.init(conn.cursor(), logSql)  # output sql.txt file
G
Ganlin Zhao 已提交
39

P
plum-lihui 已提交
40 41 42
    def checkDnodesStatusAndCreateMnode(self,dnodeNumbers):
        count=0
        while count < dnodeNumbers:
X
Xiaoyu Wang 已提交
43
            tdSql.query("select * from information_schema.ins_dnodes")
P
plum-lihui 已提交
44 45 46
            # tdLog.debug(tdSql.queryResult)
            dCnt = 0
            for i in range(dnodeNumbers):
G
Ganlin Zhao 已提交
47
                if tdSql.queryResult[i][self.dnodeStatusIndex] != "ready":
P
plum-lihui 已提交
48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65
                    break
                else:
                    dCnt += 1
            if dCnt == dnodeNumbers:
                break
            time.sleep(1)
            tdLog.debug("............... waiting for all dnodes ready!")

        tdLog.info("==============create two new mnodes ========")
        tdSql.execute("create mnode on dnode 2")
        tdSql.execute("create mnode on dnode 3")
        self.check3mnode()
        return

    def check3mnode(self):
        count=0
        while count < self.mnodeCheckCnt:
            time.sleep(1)
X
Xiaoyu Wang 已提交
66
            tdSql.query("select * from information_schema.ins_mnodes;")
G
Ganlin Zhao 已提交
67
            if tdSql.checkRows(self.mnodes) :
P
plum-lihui 已提交
68 69 70 71 72 73 74 75 76 77 78 79 80
                tdLog.debug("mnode is  three nodes")
            else:
                tdLog.exit("mnode number is correct")

            roleOfMnode0 = tdSql.queryResult[0][self.roleIndex]
            roleOfMnode1 = tdSql.queryResult[1][self.roleIndex]
            roleOfMnode2 = tdSql.queryResult[2][self.roleIndex]

            if  roleOfMnode0=='leader' and roleOfMnode1=='follower' and roleOfMnode2 == 'follower' :
                self.dnodeOfLeader = tdSql.queryResult[0][self.idIndex]
                break
            elif roleOfMnode0=='follower' and roleOfMnode1=='leader' and roleOfMnode2 == 'follower' :
                self.dnodeOfLeader = tdSql.queryResult[1][self.idIndex]
G
Ganlin Zhao 已提交
81
                break
P
plum-lihui 已提交
82 83
            elif roleOfMnode0=='follower' and roleOfMnode1=='follower' and roleOfMnode2 == 'leader' :
                self.dnodeOfLeader = tdSql.queryResult[2][self.idIndex]
G
Ganlin Zhao 已提交
84 85
                break
            else:
P
plum-lihui 已提交
86 87 88 89
                count+=1
        else:
            tdLog.exit("three mnodes is not ready in 10s ")

X
Xiaoyu Wang 已提交
90
        tdSql.query("select * from information_schema.ins_mnodes;")
G
Ganlin Zhao 已提交
91
        tdSql.checkRows(self.mnodes)
P
plum-lihui 已提交
92 93 94 95 96 97 98 99 100 101 102
        tdSql.checkData(0,self.mnodeEpIndex,'%s:%d'%(self.host,self.startPort))
        tdSql.checkData(0,self.mnodeStatusIndex,'ready')
        tdSql.checkData(1,self.mnodeEpIndex,'%s:%d'%(self.host,self.startPort+self.portStep))
        tdSql.checkData(1,self.mnodeStatusIndex,'ready')
        tdSql.checkData(2,self.mnodeEpIndex,'%s:%d'%(self.host,self.startPort+self.portStep*2))
        tdSql.checkData(2,self.mnodeStatusIndex,'ready')

    def check3mnode1off(self):
        count=0
        while count < self.mnodeCheckCnt:
            time.sleep(1)
X
Xiaoyu Wang 已提交
103
            tdSql.query("select * from information_schema.ins_mnodes")
G
Ganlin Zhao 已提交
104 105
            tdLog.debug(tdSql.queryResult)
            # if tdSql.checkRows(self.mnodes) :
P
plum-lihui 已提交
106 107 108 109 110 111 112 113 114 115 116 117 118 119
            #     tdLog.debug("mnode is three nodes")
            # else:
            #     tdLog.exit("mnode number is correct")

            roleOfMnode0 = tdSql.queryResult[0][self.roleIndex]
            roleOfMnode1 = tdSql.queryResult[1][self.roleIndex]
            roleOfMnode2 = tdSql.queryResult[2][self.roleIndex]

            if roleOfMnode0=='offline' :
                if roleOfMnode1=='leader' and roleOfMnode2 == 'follower' :
                    self.dnodeOfLeader = tdSql.queryResult[1][self.idIndex]
                    break
                elif roleOfMnode1=='follower' and roleOfMnode2 == 'leader' :
                    self.dnodeOfLeader = tdSql.queryResult[2][self.idIndex]
G
Ganlin Zhao 已提交
120
                    break
P
plum-lihui 已提交
121 122 123 124 125 126
            elif roleOfMnode1=='offline' :
                if roleOfMnode0=='leader' and roleOfMnode2 == 'follower' :
                    self.dnodeOfLeader = tdSql.queryResult[0][self.idIndex]
                    break
                elif roleOfMnode0=='follower' and roleOfMnode2 == 'leader' :
                    self.dnodeOfLeader = tdSql.queryResult[2][self.idIndex]
G
Ganlin Zhao 已提交
127
                    break
P
plum-lihui 已提交
128 129 130 131 132 133
            elif roleOfMnode2=='offline' :
                if roleOfMnode0=='leader' and roleOfMnode1 == 'follower' :
                    self.dnodeOfLeader = tdSql.queryResult[0][self.idIndex]
                    break
                elif roleOfMnode0=='follower' and roleOfMnode1 == 'leader' :
                    self.dnodeOfLeader = tdSql.queryResult[1][self.idIndex]
G
Ganlin Zhao 已提交
134
                    break
P
plum-lihui 已提交
135 136 137 138 139 140 141 142 143 144 145 146

            count+=1
        else:
            tdLog.exit("three mnodes is not ready in 10s ")

    def checkFileContent(self, consumerId, queryString):
        buildPath = tdCom.getBuildPath()
        cfgPath = tdCom.getClientCfgPath()
        dstFile = '%s/../log/dstrows_%d.txt'%(cfgPath, consumerId)
        cmdStr = '%s/build/bin/taos -c %s -s "%s >> %s"'%(buildPath, cfgPath, queryString, dstFile)
        tdLog.info(cmdStr)
        os.system(cmdStr)
G
Ganlin Zhao 已提交
147

P
plum-lihui 已提交
148 149 150 151 152
        consumeRowsFile = '%s/../log/consumerid_%d.txt'%(cfgPath, consumerId)
        tdLog.info("rows file: %s, %s"%(consumeRowsFile, dstFile))

        consumeFile = open(consumeRowsFile, mode='r')
        queryFile = open(dstFile, mode='r')
G
Ganlin Zhao 已提交
153

P
plum-lihui 已提交
154 155 156 157 158 159
        # skip first line for it is schema
        queryFile.readline()

        while True:
            dst = queryFile.readline()
            src = consumeFile.readline()
G
Ganlin Zhao 已提交
160

P
plum-lihui 已提交
161 162 163 164 165
            if dst:
                if dst != src:
                    tdLog.exit("consumerId %d consume rows is not match the rows by direct query"%consumerId)
            else:
                break
G
Ganlin Zhao 已提交
166 167
        return

P
plum-lihui 已提交
168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197
    def tmqCase1(self):
        tdLog.printNoPrefix("======== test case 1: ")
        paraDict = {'dbName':     'db1',
                    'dropFlag':   1,
                    'event':      '',
                    'vgroups':    4,
                    'stbName':    'stb',
                    'colPrefix':  'c',
                    'tagPrefix':  't',
                    'colSchema':   [{'type': 'INT', 'count':2}, {'type': 'binary', 'len':20, 'count':1}, {'type': 'TIMESTAMP', 'count':1}],
                    'tagSchema':   [{'type': 'INT', 'count':1}, {'type': 'binary', 'len':20, 'count':1}],
                    'ctbPrefix':  'ctb',
                    'ctbNum':     1,
                    'rowsPerTbl': 100000,
                    'batchNum':   10,
                    'startTs':    1640966400000,  # 2022-01-01 00:00:00.000
                    'pollDelay':  10,
                    'showMsg':    1,
                    'showRow':    1}

        topicNameList = ['topic1']
        expectRowsList = []
        tmqCom.initConsumerTable()
        tdCom.create_database(tdSql, paraDict["dbName"],paraDict["dropFlag"], vgroups=4,replica=1)
        tdLog.info("create stb")
        tdCom.create_stable(tdSql, dbname=paraDict["dbName"],stbname=paraDict["stbName"], column_elm_list=paraDict['colSchema'], tag_elm_list=paraDict['tagSchema'])
        tdLog.info("create ctb")
        tdCom.create_ctable(tdSql, dbname=paraDict["dbName"],stbname=paraDict["stbName"],tag_elm_list=paraDict['tagSchema'],count=paraDict["ctbNum"], default_ctbname_prefix=paraDict['ctbPrefix'])
        tdLog.info("async insert data")
        pThread = tmqCom.asyncInsertData(paraDict)
G
Ganlin Zhao 已提交
198

P
plum-lihui 已提交
199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237
        tdLog.info("create topics from stb with filter")
        queryString = "select ts, log(c1), ceil(pow(c1,3)) from %s.%s where c1 %% 7 == 0" %(paraDict['dbName'], paraDict['stbName'])
        sqlString = "create topic %s as %s" %(topicNameList[0], queryString)
        tdLog.info("create topic sql: %s"%sqlString)
        tdSql.execute(sqlString)

        # init consume info, and start tmq_sim, then check consume result
        tdLog.info("insert consume info to consume processor")
        consumerId   = 0
        expectrowcnt = paraDict["rowsPerTbl"] * paraDict["ctbNum"]
        topicList    = topicNameList[0]
        ifcheckdata  = 1
        ifManualCommit = 1
        keyList      = 'group.id:cgrp1, enable.auto.commit:false, auto.commit.interval.ms:6000, auto.offset.reset:earliest'
        tmqCom.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)

        tdLog.info("start consume processor")
        tmqCom.startTmqSimProcess(paraDict['pollDelay'],paraDict["dbName"],paraDict['showMsg'], paraDict['showRow'])

        tdLog.info("wait the notify info of start consume")
        tmqCom.getStartConsumeNotifyFromTmqsim()

        tdLog.info("start switch mnode ................")
        tdDnodes = cluster.dnodes

        tdLog.info("1. stop dnode 0")
        tdDnodes[0].stoptaosd()
        time.sleep(10)
        self.check3mnode1off()

        tdLog.info("2. start dnode 0")
        tdDnodes[0].starttaosd()
        self.check3mnode()

        tdLog.info("3. stop dnode 1")
        tdDnodes[1].stoptaosd()
        time.sleep(10)
        self.check3mnode1off()

G
Ganlin Zhao 已提交
238 239 240 241 242
        tdLog.info("switch end and wait insert data end ................")
        pThread.join()

        tdLog.info("check the consume result")
        tdSql.query(queryString)
P
plum-lihui 已提交
243 244 245 246
        expectRowsList.append(tdSql.getRows())

        expectRows = 1
        resultList = tmqCom.selectConsumeResult(expectRows)
G
Ganlin Zhao 已提交
247

P
plum-lihui 已提交
248 249 250 251
        if expectRowsList[0] != resultList[0]:
            tdLog.info("expect consume rows: %d, act consume rows: %d"%(expectRowsList[0], resultList[0]))
            tdLog.exit("0 tmq consume rows error!")

G
Ganlin Zhao 已提交
252
        self.checkFileContent(consumerId, queryString)
P
plum-lihui 已提交
253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272

        time.sleep(10)
        for i in range(len(topicNameList)):
            tdSql.query("drop topic %s"%topicNameList[i])

        tdLog.printNoPrefix("======== test case 1 end ...... ")

    def run(self):
        tdLog.printNoPrefix("======== Notes: must add '-N 5' for run the script ========")
        self.checkDnodesStatusAndCreateMnode(self.dnodes)
        self.tmqCase1()

    def stop(self):
        tdSql.close()
        tdLog.success(f"{__file__} successfully executed")

event = threading.Event()

tdCases.addLinux(__file__, TDTestCase())
tdCases.addWindows(__file__, TDTestCase())