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

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)

23
    def init(self, conn, logSql, replicaVar=1):
P
plum-lihui 已提交
24
        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
                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

G
Ganlin Zhao 已提交
52
    def initConsumerTable(self,cdbName='cdb'):
P
plum-lihui 已提交
53 54 55 56 57 58 59 60
        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)

G
Ganlin Zhao 已提交
61
    def insertConsumerInfo(self,consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifmanualcommit,cdbName='cdb'):
P
plum-lihui 已提交
62 63 64 65 66 67 68 69 70 71 72 73 74
        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
                time.sleep(5)
G
Ganlin Zhao 已提交
76

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
            resultList.append(tdSql.getData(i , 3))
G
Ganlin Zhao 已提交
80

P
plum-lihui 已提交
81 82 83 84 85 86 87
        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 '
G
Ganlin Zhao 已提交
88

wafwerar's avatar
wafwerar 已提交
89 90
        if (platform.system().lower() == 'windows'):
            shellCmd = 'mintty -h never -w hide ' + buildPath + '\\build\\bin\\tmq_sim.exe -c ' + cfgPath
G
Ganlin Zhao 已提交
91 92
            shellCmd += " -y %d -d %s -g %d -r %d -w %s "%(pollDelay, dbName, showMsg, showRow, cdbName)
            shellCmd += "> nul 2>&1 &"
wafwerar's avatar
wafwerar 已提交
93 94
        else:
            shellCmd = 'nohup ' + buildPath + '/build/bin/tmq_sim -c ' + cfgPath
G
Ganlin Zhao 已提交
95
            shellCmd += " -y %d -d %s -g %d -r %d -w %s "%(pollDelay, dbName, showMsg, showRow, cdbName)
wafwerar's avatar
wafwerar 已提交
96
            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):
G
Ganlin Zhao 已提交
101
        tsql.execute("create database if not exists %s vgroups %d"%(dbName, vgroups))
P
plum-lihui 已提交
102 103 104 105 106 107 108 109 110 111 112 113
        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)
G
Ganlin Zhao 已提交
114 115

        event.set()
P
plum-lihui 已提交
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
        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
G
Ganlin Zhao 已提交
144

P
plum-lihui 已提交
145 146 147 148 149 150 151 152 153
    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

        self.insert_data(tsql,\
                         parameterDict["dbName"],\
                         parameterDict["stbName"],\
                         parameterDict["ctbNum"],\
                         parameterDict["rowsPerTbl"],\
                         parameterDict["batchNum"],\
G
Ganlin Zhao 已提交
162
                         parameterDict["startTs"])
P
plum-lihui 已提交
163 164
        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
        topicName1 = 'topic_db60'
        topicName2 = 'topic_db61'
G
Ganlin Zhao 已提交
204 205

        tdSql.execute("create topic %s as database %s" %(topicName1, parameterDict['dbName']))
P
plum-lihui 已提交
206
        tdSql.execute("create topic %s as database %s" %(topicName2, parameterDict2['dbName']))
G
Ganlin Zhao 已提交
207

P
plum-lihui 已提交
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
        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)
G
Ganlin Zhao 已提交
218

P
plum-lihui 已提交
219 220
        #consumerId   = 1
        #self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
G
Ganlin Zhao 已提交
221

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
G
Ganlin Zhao 已提交
231
        prepareEnvThread.join()
P
plum-lihui 已提交
232
        prepareEnvThread2.join()
G
Ganlin Zhao 已提交
233

P
plum-lihui 已提交
234 235 236 237 238 239
        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]
G
Ganlin Zhao 已提交
240

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
        topicName1 = 'topic_db60'
        topicName2 = 'topic_db61'
G
Ganlin Zhao 已提交
289 290

        tdSql.execute("create topic %s as database %s" %(topicName1, parameterDict['dbName']))
P
plum-lihui 已提交
291
        tdSql.execute("create topic %s as database %s" %(topicName2, parameterDict2['dbName']))
G
Ganlin Zhao 已提交
292

P
plum-lihui 已提交
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
                        auto.offset.reset:earliest'
        self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
G
Ganlin Zhao 已提交
303

P
plum-lihui 已提交
304 305
        consumerId   = 1
        self.insertConsumerInfo(consumerId, expectrowcnt,topicList,keyList,ifcheckdata,ifManualCommit)
G
Ganlin Zhao 已提交
306

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
        showMsg   = 1
        showRow   = 1
        self.startTmqSimProcess(buildPath,cfgPath,pollDelay,parameterDict["dbName"],showMsg, showRow)

        # wait for data ready
G
Ganlin Zhao 已提交
316
        prepareEnvThread.join()
P
plum-lihui 已提交
317
        prepareEnvThread2.join()
G
Ganlin Zhao 已提交
318

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
        resultList = self.selectConsumeResult(expectRows)
        totalConsumeRows = 0
        for i in range(expectRows):
            totalConsumeRows += resultList[i]
G
Ganlin Zhao 已提交
325

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())