subscribeDb1.py 15.6 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

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 *

class TDTestCase:
    hostname = socket.gethostname()
    #rpcDebugFlagVal = '143'
    #clientCfgDict = {'serverPort': '', 'firstEp': '', 'secondEp':'', 'rpcDebugFlag':'135', 'fqdn':''}
    #clientCfgDict["rpcDebugFlag"]  = rpcDebugFlagVal
    #updatecfgDict = {'clientCfg': {}, 'serverPort': '', 'firstEp': '', 'secondEp':'', 'rpcDebugFlag':'135', 'fqdn':''}
    #updatecfgDict["rpcDebugFlag"] = rpcDebugFlagVal
    #print ("===================: ", updatecfgDict)

    def init(self, conn, logSql):
        tdLog.debug(f"start to excute {__file__}")
25 26
        tdSql.init(conn.cursor())
        #tdSql.init(conn.cursor(), logSql)  # output sql.txt file
P
plum-lihui 已提交
27 28 29 30 31 32 33 34 35 36

    def getBuildPath(self):
        selfPath = os.path.dirname(os.path.realpath(__file__))

        if ("community" in selfPath):
            projPath = selfPath[:selfPath.find("community")]
        else:
            projPath = selfPath[:selfPath.find("tests")]

        for root, dirs, files in os.walk(projPath):
wafwerar's avatar
wafwerar 已提交
37
            if ("taosd" in files or "taosd.exe" in files):
P
plum-lihui 已提交
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
                rootRealPath = os.path.dirname(os.path.realpath(root))
                if ("packaging" not in rootRealPath):
                    buildPath = root[:len(root) - len("/build/bin")]
                    break
        return buildPath

    def newcur(self,cfg,host,port):
        user = "root"
        password = "taosdata"
        con=taos.connect(host=host, user=user, password=password, config=cfg ,port=port)
        cur=con.cursor()
        print(cur)
        return cur

    def initConsumerTable(self,cdbName='cdb'):        
        tdLog.info("create consume database, and consume info table, and consume result table")
        tdSql.query("create database if not exists %s vgroups 1"%(cdbName))
        tdSql.query("drop table if exists %s.consumeinfo "%(cdbName))
        tdSql.query("drop table if exists %s.consumeresult "%(cdbName))

        tdSql.query("create table %s.consumeinfo (ts timestamp, consumerid int, topiclist binary(1024), keylist binary(1024), expectmsgcnt bigint, ifcheckdata int, ifmanualcommit int)"%cdbName)
        tdSql.query("create table %s.consumeresult (ts timestamp, consumerid int, consummsgcnt bigint, consumrowcnt bigint, checkresult int)"%cdbName)

    def insertConsumerInfo(self,consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifmanualcommit,cdbName='cdb'):    
        sql = "insert into %s.consumeinfo values "%cdbName
        sql += "(now, %d, '%s', '%s', %d, %d, %d)"%(consumerId, topicList, keyList, expectrowcnt, ifcheckdata, ifmanualcommit)
        tdLog.info("consume info sql: %s"%sql)
        tdSql.query(sql)

    def selectConsumeResult(self,expectRows,cdbName='cdb'):
        resultList=[]
        while 1:
            tdSql.query("select * from %s.consumeresult"%cdbName)
            #tdLog.info("row: %d, %l64d, %l64d"%(tdSql.getData(0, 1),tdSql.getData(0, 2),tdSql.getData(0, 3))
            if tdSql.getRows() == expectRows:
                break
            else:
P
plum-lihui 已提交
75 76
                time.sleep(5)
        
P
plum-lihui 已提交
77
        for i in range(expectRows):
P
plum-lihui 已提交
78
            tdLog.info ("consume id: %d, consume msgs: %d, consume rows: %d"%(tdSql.getData(i , 1), tdSql.getData(i , 2), tdSql.getData(i , 3)))
P
plum-lihui 已提交
79 80 81 82 83 84 85 86 87
            resultList.append(tdSql.getData(i , 3))
        
        return resultList

    def startTmqSimProcess(self,buildPath,cfgPath,pollDelay,dbName,showMsg=1,showRow=1,cdbName='cdb',valgrind=0):
        if valgrind == 1:
            logFile = cfgPath + '/../log/valgrind-tmq.log'
            shellCmd = 'nohup valgrind --log-file=' + logFile
            shellCmd += '--tool=memcheck --leak-check=full --show-reachable=no --track-origins=yes --show-leak-kinds=all --num-callers=20 -v --workaround-gcc296-bugs=yes '
P
plum-lihui 已提交
88
        
wafwerar's avatar
wafwerar 已提交
89 90 91 92 93 94 95 96
        if (platform.system().lower() == 'windows'):
            shellCmd = 'mintty -h never -w hide ' + buildPath + '\\build\\bin\\tmq_sim.exe -c ' + cfgPath
            shellCmd += " -y %d -d %s -g %d -r %d -w %s "%(pollDelay, dbName, showMsg, showRow, cdbName) 
            shellCmd += "> nul 2>&1 &"   
        else:
            shellCmd = 'nohup ' + buildPath + '/build/bin/tmq_sim -c ' + cfgPath
            shellCmd += " -y %d -d %s -g %d -r %d -w %s "%(pollDelay, dbName, showMsg, showRow, cdbName) 
            shellCmd += "> /dev/null 2>&1 &"
P
plum-lihui 已提交
97 98 99
        tdLog.info(shellCmd)
        os.system(shellCmd)

P
plum-lihui 已提交
100
    def create_tables(self,tsql, dbName,vgroups,stbName,ctbNum):
P
plum-lihui 已提交
101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153
        tsql.execute("create database if not exists %s vgroups %d"%(dbName, vgroups))        
        tsql.execute("use %s" %dbName)
        tsql.execute("create table  if not exists %s (ts timestamp, c1 bigint, c2 binary(16)) tags(t1 int)"%stbName)
        pre_create = "create table"
        sql = pre_create
        #tdLog.debug("doing create one  stable %s and %d  child table in %s  ..." %(stbname, count ,dbname))
        for i in range(ctbNum):
            sql += " %s_%d using %s tags(%d)"%(stbName,i,stbName,i+1)
            if (i > 0) and (i%100 == 0):
                tsql.execute(sql)
                sql = pre_create
        if sql != pre_create:
            tsql.execute(sql)
        
        event.set()    
        tdLog.debug("complete to create database[%s], stable[%s] and %d child tables" %(dbName, stbName, ctbNum))
        return

    def insert_data(self,tsql,dbName,stbName,ctbNum,rowsPerTbl,batchNum,startTs):
        tdLog.debug("start to insert data ............")
        tsql.execute("use %s" %dbName)
        pre_insert = "insert into "
        sql = pre_insert

        t = time.time()
        startTs = int(round(t * 1000))
        #tdLog.debug("doing insert data into stable:%s rows:%d ..."%(stbName, allRows))
        for i in range(ctbNum):
            sql += " %s_%d values "%(stbName,i)
            for j in range(rowsPerTbl):
                sql += "(%d, %d, 'tmqrow_%d') "%(startTs + j, j, j)
                if (j > 0) and ((j%batchNum == 0) or (j == rowsPerTbl - 1)):
                    tsql.execute(sql)
                    if j < rowsPerTbl - 1:
                        sql = "insert into %s_%d values " %(stbName,i)
                    else:
                        sql = "insert into "
        #end sql
        if sql != pre_insert:
            #print("insert sql:%s"%sql)
            tsql.execute(sql)
        tdLog.debug("insert data ............ [OK]")
        return
        
    def prepareEnv(self, **parameterDict):
        print ("input parameters:")
        print (parameterDict)
        # create new connector for my thread
        tsql=self.newcur(parameterDict['cfg'], 'localhost', 6030)
        self.create_tables(tsql,\
                           parameterDict["dbName"],\
                           parameterDict["vgroups"],\
                           parameterDict["stbName"],\
P
plum-lihui 已提交
154
                           parameterDict["ctbNum"])
P
plum-lihui 已提交
155 156 157 158 159 160 161 162 163 164

        self.insert_data(tsql,\
                         parameterDict["dbName"],\
                         parameterDict["stbName"],\
                         parameterDict["ctbNum"],\
                         parameterDict["rowsPerTbl"],\
                         parameterDict["batchNum"],\
                         parameterDict["startTs"])                           
        return

P
plum-lihui 已提交
165 166
    def tmqCase6(self, cfgPath, buildPath):
        tdLog.printNoPrefix("======== test case 6: Produce while one consumers to subscribe tow topic, Each contains one db")
P
plum-lihui 已提交
167 168 169
        tdLog.info("step 1: create database, stb, ctb and insert data")
        # create and start thread
        parameterDict = {'cfg':        '',       \
P
plum-lihui 已提交
170
                         'dbName':     'db60',    \
P
plum-lihui 已提交
171 172 173
                         'vgroups':    4,        \
                         'stbName':    'stb',    \
                         'ctbNum':     10,       \
P
plum-lihui 已提交
174
                         'rowsPerTbl': 5000,    \
P
plum-lihui 已提交
175 176 177 178 179 180 181 182 183 184 185
                         'batchNum':   100,      \
                         'startTs':    1640966400000}  # 2022-01-01 00:00:00.000
        parameterDict['cfg'] = cfgPath

        self.initConsumerTable()

        tdSql.execute("create database if not exists %s vgroups %d" %(parameterDict['dbName'], parameterDict['vgroups']))

        prepareEnvThread = threading.Thread(target=self.prepareEnv, kwargs=parameterDict)
        prepareEnvThread.start()

P
plum-lihui 已提交
186 187
        parameterDict2 = {'cfg':        '',       \
                         'dbName':     'db61',    \
P
plum-lihui 已提交
188
                         'vgroups':    4,        \
P
plum-lihui 已提交
189
                         'stbName':    'stb2',    \
P
plum-lihui 已提交
190
                         'ctbNum':     10,       \
P
plum-lihui 已提交
191
                         'rowsPerTbl': 5000,    \
P
plum-lihui 已提交
192 193 194 195
                         'batchNum':   100,      \
                         'startTs':    1640966400000}  # 2022-01-01 00:00:00.000
        parameterDict['cfg'] = cfgPath

P
plum-lihui 已提交
196
        tdSql.execute("create database if not exists %s vgroups %d" %(parameterDict2['dbName'], parameterDict2['vgroups']))
P
plum-lihui 已提交
197

P
plum-lihui 已提交
198 199
        prepareEnvThread2 = threading.Thread(target=self.prepareEnv, kwargs=parameterDict2)
        prepareEnvThread2.start()
P
plum-lihui 已提交
200 201

        tdLog.info("create topics from db")
P
plum-lihui 已提交
202 203 204 205 206
        topicName1 = 'topic_db60'
        topicName2 = 'topic_db61'
        
        tdSql.execute("create topic %s as database %s" %(topicName1, parameterDict['dbName']))        
        tdSql.execute("create topic %s as database %s" %(topicName2, parameterDict2['dbName']))
P
plum-lihui 已提交
207 208
        
        consumerId   = 0
P
plum-lihui 已提交
209 210
        expectrowcnt = parameterDict["rowsPerTbl"] * parameterDict["ctbNum"] +  parameterDict2["rowsPerTbl"] * parameterDict2["ctbNum"]
        topicList    = topicName1 + ',' + topicName2
P
plum-lihui 已提交
211
        ifcheckdata  = 0
P
plum-lihui 已提交
212
        ifManualCommit = 0
P
plum-lihui 已提交
213 214 215 216 217 218
        keyList      = 'group.id:cgrp1,\
                        enable.auto.commit:false,\
                        auto.commit.interval.ms:6000,\
                        auto.offset.reset:earliest'
        self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
        
P
plum-lihui 已提交
219 220 221
        #consumerId   = 1
        #self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
        
P
plum-lihui 已提交
222 223 224
        event.wait()

        tdLog.info("start consume processor")
225
        pollDelay = 100
P
plum-lihui 已提交
226 227 228 229 230
        showMsg   = 1
        showRow   = 1
        self.startTmqSimProcess(buildPath,cfgPath,pollDelay,parameterDict["dbName"],showMsg, showRow)

        # wait for data ready
P
plum-lihui 已提交
231 232
        prepareEnvThread.join()        
        prepareEnvThread2.join()
P
plum-lihui 已提交
233 234 235 236 237 238 239 240
        
        tdLog.info("insert process end, and start to check consume result")
        expectRows = 1
        resultList = self.selectConsumeResult(expectRows)
        totalConsumeRows = 0
        for i in range(expectRows):
            totalConsumeRows += resultList[i]
        
P
plum-lihui 已提交
241 242
        if totalConsumeRows != expectrowcnt:
            tdLog.info("act consume rows: %d, expect consume rows: %d"%(totalConsumeRows, expectrowcnt))
P
plum-lihui 已提交
243 244 245
            tdLog.exit("tmq consume rows error!")

        tdSql.query("drop topic %s"%topicName1)
P
plum-lihui 已提交
246
        tdSql.query("drop topic %s"%topicName2)
P
plum-lihui 已提交
247

P
plum-lihui 已提交
248
        tdLog.printNoPrefix("======== test case 6 end ...... ")
P
plum-lihui 已提交
249

P
plum-lihui 已提交
250 251
    def tmqCase7(self, cfgPath, buildPath):
        tdLog.printNoPrefix("======== test case 7: Produce while two consumers to subscribe tow topic, Each contains one db")
P
plum-lihui 已提交
252 253 254
        tdLog.info("step 1: create database, stb, ctb and insert data")
        # create and start thread
        parameterDict = {'cfg':        '',       \
P
plum-lihui 已提交
255
                         'dbName':     'db70',    \
P
plum-lihui 已提交
256 257 258
                         'vgroups':    4,        \
                         'stbName':    'stb',    \
                         'ctbNum':     10,       \
P
plum-lihui 已提交
259
                         'rowsPerTbl': 5000,    \
P
plum-lihui 已提交
260 261 262 263 264 265 266 267 268 269 270
                         'batchNum':   100,      \
                         'startTs':    1640966400000}  # 2022-01-01 00:00:00.000
        parameterDict['cfg'] = cfgPath

        self.initConsumerTable()

        tdSql.execute("create database if not exists %s vgroups %d" %(parameterDict['dbName'], parameterDict['vgroups']))

        prepareEnvThread = threading.Thread(target=self.prepareEnv, kwargs=parameterDict)
        prepareEnvThread.start()

P
plum-lihui 已提交
271 272
        parameterDict2 = {'cfg':        '',       \
                         'dbName':     'db71',    \
P
plum-lihui 已提交
273
                         'vgroups':    4,        \
P
plum-lihui 已提交
274
                         'stbName':    'stb2',    \
P
plum-lihui 已提交
275
                         'ctbNum':     10,       \
P
plum-lihui 已提交
276
                         'rowsPerTbl': 5000,    \
P
plum-lihui 已提交
277 278 279 280
                         'batchNum':   100,      \
                         'startTs':    1640966400000}  # 2022-01-01 00:00:00.000
        parameterDict['cfg'] = cfgPath

P
plum-lihui 已提交
281
        tdSql.execute("create database if not exists %s vgroups %d" %(parameterDict2['dbName'], parameterDict2['vgroups']))
P
plum-lihui 已提交
282

P
plum-lihui 已提交
283 284
        prepareEnvThread2 = threading.Thread(target=self.prepareEnv, kwargs=parameterDict2)
        prepareEnvThread2.start()
P
plum-lihui 已提交
285 286

        tdLog.info("create topics from db")
P
plum-lihui 已提交
287 288 289 290 291
        topicName1 = 'topic_db60'
        topicName2 = 'topic_db61'
        
        tdSql.execute("create topic %s as database %s" %(topicName1, parameterDict['dbName']))        
        tdSql.execute("create topic %s as database %s" %(topicName2, parameterDict2['dbName']))
P
plum-lihui 已提交
292 293
        
        consumerId   = 0
P
plum-lihui 已提交
294 295
        expectrowcnt = parameterDict["rowsPerTbl"] * parameterDict["ctbNum"] +  parameterDict2["rowsPerTbl"] * parameterDict2["ctbNum"]
        topicList    = topicName1 + ',' + topicName2
P
plum-lihui 已提交
296 297 298
        ifcheckdata  = 0
        ifManualCommit = 1
        keyList      = 'group.id:cgrp1,\
P
plum-lihui 已提交
299 300
                        enable.auto.commit:false,\
                        auto.commit.interval.ms:6000,\
P
plum-lihui 已提交
301 302 303
                        auto.offset.reset:earliest'
        self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
        
P
plum-lihui 已提交
304 305 306
        consumerId   = 1
        self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
        
P
plum-lihui 已提交
307 308 309
        event.wait()

        tdLog.info("start consume processor")
P
plum-lihui 已提交
310
        pollDelay = 100
P
plum-lihui 已提交
311 312 313 314 315 316
        showMsg   = 1
        showRow   = 1
        self.startTmqSimProcess(buildPath,cfgPath,pollDelay,parameterDict["dbName"],showMsg, showRow)

        # wait for data ready
        prepareEnvThread.join()        
P
plum-lihui 已提交
317 318
        prepareEnvThread2.join()
        
P
plum-lihui 已提交
319
        tdLog.info("insert process end, and start to check consume result")
P
plum-lihui 已提交
320
        expectRows = 2
P
plum-lihui 已提交
321 322 323 324 325
        resultList = self.selectConsumeResult(expectRows)
        totalConsumeRows = 0
        for i in range(expectRows):
            totalConsumeRows += resultList[i]
        
P
plum-lihui 已提交
326 327
        if totalConsumeRows != expectrowcnt:
            tdLog.info("act consume rows: %d, expect consume rows: %d"%(totalConsumeRows, expectrowcnt))
P
plum-lihui 已提交
328 329 330
            tdLog.exit("tmq consume rows error!")

        tdSql.query("drop topic %s"%topicName1)
P
plum-lihui 已提交
331
        tdSql.query("drop topic %s"%topicName2)
P
plum-lihui 已提交
332

P
plum-lihui 已提交
333
        tdLog.printNoPrefix("======== test case 7 end ...... ")
P
plum-lihui 已提交
334 335 336 337 338 339 340 341 342 343 344 345

    def run(self):
        tdSql.prepare()

        buildPath = self.getBuildPath()
        if (buildPath == ""):
            tdLog.exit("taosd not found!")
        else:
            tdLog.info("taosd found in %s" % buildPath)
        cfgPath = buildPath + "/../sim/psim/cfg"
        tdLog.info("cfgPath: %s" % cfgPath)

P
plum-lihui 已提交
346 347 348
        self.tmqCase6(cfgPath, buildPath)
        self.tmqCase7(cfgPath, buildPath)

P
plum-lihui 已提交
349 350 351 352 353 354 355 356 357

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

event = threading.Event()

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