parent
bd79f4c26b
commit
d4fdf17cae
File diff suppressed because it is too large
Load Diff
|
@ -12,12 +12,13 @@
|
||||||
# -*- coding: utf-8 -*-
|
# -*- coding: utf-8 -*-
|
||||||
|
|
||||||
import taos
|
import taos
|
||||||
from util.log import *
|
|
||||||
from util.cases import *
|
|
||||||
from util.sql import *
|
|
||||||
from util.common import *
|
|
||||||
from util.sqlset import *
|
|
||||||
from taos.tmq import *
|
from taos.tmq import *
|
||||||
|
from util.cases import *
|
||||||
|
from util.common import *
|
||||||
|
from util.log import *
|
||||||
|
from util.sql import *
|
||||||
|
from util.sqlset import *
|
||||||
|
|
||||||
|
|
||||||
class TDTestCase:
|
class TDTestCase:
|
||||||
def init(self, conn, logSql, replicaVar=1):
|
def init(self, conn, logSql, replicaVar=1):
|
||||||
|
@ -26,10 +27,10 @@ class TDTestCase:
|
||||||
tdSql.init(conn.cursor())
|
tdSql.init(conn.cursor())
|
||||||
self.setsql = TDSetSql()
|
self.setsql = TDSetSql()
|
||||||
self.stbname = 'stb'
|
self.stbname = 'stb'
|
||||||
self.binary_length = 20 # the length of binary for column_dict
|
self.binary_length = 20 # the length of binary for column_dict
|
||||||
self.nchar_length = 20 # the length of nchar for column_dict
|
self.nchar_length = 20 # the length of nchar for column_dict
|
||||||
self.column_dict = {
|
self.column_dict = {
|
||||||
'ts' : 'timestamp',
|
'ts': 'timestamp',
|
||||||
'col1': 'tinyint',
|
'col1': 'tinyint',
|
||||||
'col2': 'smallint',
|
'col2': 'smallint',
|
||||||
'col3': 'int',
|
'col3': 'int',
|
||||||
|
@ -45,7 +46,7 @@ class TDTestCase:
|
||||||
'col13': f'nchar({self.nchar_length})'
|
'col13': f'nchar({self.nchar_length})'
|
||||||
}
|
}
|
||||||
self.tag_dict = {
|
self.tag_dict = {
|
||||||
'ts_tag' : 'timestamp',
|
'ts_tag': 'timestamp',
|
||||||
't1': 'tinyint',
|
't1': 'tinyint',
|
||||||
't2': 'smallint',
|
't2': 'smallint',
|
||||||
't3': 'int',
|
't3': 'int',
|
||||||
|
@ -67,25 +68,28 @@ class TDTestCase:
|
||||||
f'now,1,2,3,4,5,6,7,8,9.9,10.1,true,"abcd","涛思数据"'
|
f'now,1,2,3,4,5,6,7,8,9.9,10.1,true,"abcd","涛思数据"'
|
||||||
]
|
]
|
||||||
self.tbnum = 1
|
self.tbnum = 1
|
||||||
|
|
||||||
def prepare_data(self):
|
def prepare_data(self):
|
||||||
tdSql.execute(self.setsql.set_create_stable_sql(self.stbname,self.column_dict,self.tag_dict))
|
tdSql.execute(self.setsql.set_create_stable_sql(self.stbname, self.column_dict, self.tag_dict))
|
||||||
for i in range(self.tbnum):
|
for i in range(self.tbnum):
|
||||||
tdSql.execute(f'create table {self.stbname}_{i} using {self.stbname} tags({self.tag_list[i]})')
|
tdSql.execute(f'create table {self.stbname}_{i} using {self.stbname} tags({self.tag_list[i]})')
|
||||||
for j in self.values_list:
|
for j in self.values_list:
|
||||||
tdSql.execute(f'insert into {self.stbname}_{i} values({j})')
|
tdSql.execute(f'insert into {self.stbname}_{i} values({j})')
|
||||||
|
|
||||||
def create_user(self):
|
def create_user(self):
|
||||||
for user_name in ['jiacy1_all','jiacy1_read','jiacy1_write','jiacy1_none','jiacy0_all','jiacy0_read','jiacy0_write','jiacy0_none']:
|
for user_name in ['jiacy1_all', 'jiacy1_read', 'jiacy1_write', 'jiacy1_none', 'jiacy0_all', 'jiacy0_read',
|
||||||
|
'jiacy0_write', 'jiacy0_none']:
|
||||||
if 'jiacy1' in user_name.lower():
|
if 'jiacy1' in user_name.lower():
|
||||||
tdSql.execute(f'create user {user_name} pass "123" sysinfo 1')
|
tdSql.execute(f'create user {user_name} pass "123" sysinfo 1')
|
||||||
elif 'jiacy0' in user_name.lower():
|
elif 'jiacy0' in user_name.lower():
|
||||||
tdSql.execute(f'create user {user_name} pass "123" sysinfo 0')
|
tdSql.execute(f'create user {user_name} pass "123" sysinfo 0')
|
||||||
for user_name in ['jiacy1_all','jiacy1_read','jiacy0_all','jiacy0_read']:
|
for user_name in ['jiacy1_all', 'jiacy1_read', 'jiacy0_all', 'jiacy0_read']:
|
||||||
tdSql.execute(f'grant read on db to {user_name}')
|
tdSql.execute(f'grant read on db to {user_name}')
|
||||||
for user_name in ['jiacy1_all','jiacy1_write','jiacy0_all','jiacy0_write']:
|
for user_name in ['jiacy1_all', 'jiacy1_write', 'jiacy0_all', 'jiacy0_write']:
|
||||||
tdSql.execute(f'grant write on db to {user_name}')
|
tdSql.execute(f'grant write on db to {user_name}')
|
||||||
|
|
||||||
def user_privilege_check(self):
|
def user_privilege_check(self):
|
||||||
jiacy1_read_conn = taos.connect(user='jiacy1_read',password='123')
|
jiacy1_read_conn = taos.connect(user='jiacy1_read', password='123')
|
||||||
sql = "create table ntb (ts timestamp,c0 int)"
|
sql = "create table ntb (ts timestamp,c0 int)"
|
||||||
expectErrNotOccured = True
|
expectErrNotOccured = True
|
||||||
try:
|
try:
|
||||||
|
@ -94,32 +98,34 @@ class TDTestCase:
|
||||||
expectErrNotOccured = False
|
expectErrNotOccured = False
|
||||||
if expectErrNotOccured:
|
if expectErrNotOccured:
|
||||||
caller = inspect.getframeinfo(inspect.stack()[1][0])
|
caller = inspect.getframeinfo(inspect.stack()[1][0])
|
||||||
tdLog.exit(f"{caller.filename}({caller.lineno}) failed: sql:{sql}, expect error not occured" )
|
tdLog.exit(f"{caller.filename}({caller.lineno}) failed: sql:{sql}, expect error not occured")
|
||||||
else:
|
else:
|
||||||
self.queryRows = 0
|
self.queryRows = 0
|
||||||
self.queryCols = 0
|
self.queryCols = 0
|
||||||
self.queryResult = None
|
self.queryResult = None
|
||||||
tdLog.info(f"sql:{sql}, expect error occured")
|
tdLog.info(f"sql:{sql}, expect error occured")
|
||||||
pass
|
pass
|
||||||
|
|
||||||
def drop_topic(self):
|
def drop_topic(self):
|
||||||
jiacy1_all_conn = taos.connect(user='jiacy1_all',password='123')
|
jiacy1_all_conn = taos.connect(user='jiacy1_all', password='123')
|
||||||
jiacy1_read_conn = taos.connect(user='jiacy1_read',password='123')
|
jiacy1_read_conn = taos.connect(user='jiacy1_read', password='123')
|
||||||
jiacy1_write_conn = taos.connect(user='jiacy1_write',password='123')
|
jiacy1_write_conn = taos.connect(user='jiacy1_write', password='123')
|
||||||
jiacy1_none_conn = taos.connect(user='jiacy1_none',password='123')
|
jiacy1_none_conn = taos.connect(user='jiacy1_none', password='123')
|
||||||
jiacy0_all_conn = taos.connect(user='jiacy0_all',password='123')
|
jiacy0_all_conn = taos.connect(user='jiacy0_all', password='123')
|
||||||
jiacy0_read_conn = taos.connect(user='jiacy0_read',password='123')
|
jiacy0_read_conn = taos.connect(user='jiacy0_read', password='123')
|
||||||
jiacy0_write_conn = taos.connect(user='jiacy0_write',password='123')
|
jiacy0_write_conn = taos.connect(user='jiacy0_write', password='123')
|
||||||
jiacy0_none_conn = taos.connect(user='jiacy0_none',password='123')
|
jiacy0_none_conn = taos.connect(user='jiacy0_none', password='123')
|
||||||
tdSql.execute('create topic root_db as select * from db.stb')
|
tdSql.execute('create topic root_db as select * from db.stb')
|
||||||
for user in [jiacy1_all_conn,jiacy1_read_conn,jiacy0_all_conn,jiacy0_read_conn]:
|
for user in [jiacy1_all_conn, jiacy1_read_conn, jiacy0_all_conn, jiacy0_read_conn]:
|
||||||
user.execute(f'create topic db_jiacy as select * from db.stb')
|
user.execute(f'create topic db_jiacy as select * from db.stb')
|
||||||
user.execute('drop topic db_jiacy')
|
user.execute('drop topic db_jiacy')
|
||||||
for user in [jiacy1_write_conn,jiacy1_none_conn,jiacy0_write_conn,jiacy0_none_conn,jiacy1_all_conn,jiacy1_read_conn,jiacy0_all_conn,jiacy0_read_conn]:
|
for user in [jiacy1_write_conn, jiacy1_none_conn, jiacy0_write_conn, jiacy0_none_conn, jiacy1_all_conn,
|
||||||
|
jiacy1_read_conn, jiacy0_all_conn, jiacy0_read_conn]:
|
||||||
sql_list = []
|
sql_list = []
|
||||||
if user in [jiacy1_all_conn,jiacy1_read_conn,jiacy0_all_conn,jiacy0_read_conn]:
|
if user in [jiacy1_all_conn, jiacy1_read_conn, jiacy0_all_conn, jiacy0_read_conn]:
|
||||||
sql_list = ['drop topic root_db']
|
sql_list = ['drop topic root_db']
|
||||||
elif user in [jiacy1_write_conn,jiacy1_none_conn,jiacy0_write_conn,jiacy0_none_conn]:
|
elif user in [jiacy1_write_conn, jiacy1_none_conn, jiacy0_write_conn, jiacy0_none_conn]:
|
||||||
sql_list = ['drop topic root_db','create topic db_jiacy as select * from db.stb']
|
sql_list = ['drop topic root_db', 'create topic db_jiacy as select * from db.stb']
|
||||||
for sql in sql_list:
|
for sql in sql_list:
|
||||||
expectErrNotOccured = True
|
expectErrNotOccured = True
|
||||||
try:
|
try:
|
||||||
|
@ -128,33 +134,26 @@ class TDTestCase:
|
||||||
expectErrNotOccured = False
|
expectErrNotOccured = False
|
||||||
if expectErrNotOccured:
|
if expectErrNotOccured:
|
||||||
caller = inspect.getframeinfo(inspect.stack()[1][0])
|
caller = inspect.getframeinfo(inspect.stack()[1][0])
|
||||||
tdLog.exit(f"{caller.filename}({caller.lineno}) failed: sql:{sql}, expect error not occured" )
|
tdLog.exit(f"{caller.filename}({caller.lineno}) failed: sql:{sql}, expect error not occured")
|
||||||
else:
|
else:
|
||||||
self.queryRows = 0
|
self.queryRows = 0
|
||||||
self.queryCols = 0
|
self.queryCols = 0
|
||||||
self.queryResult = None
|
self.queryResult = None
|
||||||
tdLog.info(f"sql:{sql}, expect error occured")
|
tdLog.info(f"sql:{sql}, expect error occured")
|
||||||
|
|
||||||
def tmq_commit_cb_print(tmq, resp, param=None):
|
def tmq_commit_cb_print(tmq, resp, param=None):
|
||||||
print(f"commit: {resp}, tmq: {tmq}, param: {param}")
|
print(f"commit: {resp}, tmq: {tmq}, param: {param}")
|
||||||
|
|
||||||
def subscribe_topic(self):
|
def subscribe_topic(self):
|
||||||
print("create topic")
|
print("create topic")
|
||||||
tdSql.execute('create topic db_topic as select * from db.stb')
|
tdSql.execute('create topic db_topic as select * from db.stb')
|
||||||
tdSql.execute('grant subscribe on db_topic to jiacy1_all')
|
tdSql.execute('grant subscribe on db_topic to jiacy1_all')
|
||||||
print("build consumer")
|
print("build consumer")
|
||||||
conf = TaosTmqConf()
|
tmq = Consumer({"group.id": "tg2", "td.connect.user": "jiacy1_all", "td.connect.pass": "123",
|
||||||
conf.set("group.id", "tg2")
|
"enable.auto.commit": "true"})
|
||||||
conf.set("td.connect.user", "jiacy1_all")
|
|
||||||
conf.set("td.connect.pass", "123")
|
|
||||||
conf.set("enable.auto.commit", "true")
|
|
||||||
conf.set_auto_commit_cb(self.tmq_commit_cb_print, None)
|
|
||||||
tmq = conf.new_consumer()
|
|
||||||
print("build topic list")
|
print("build topic list")
|
||||||
topic_list = TaosTmqList()
|
tmq.subscribe(["db_topic"])
|
||||||
topic_list.append("db_topic")
|
|
||||||
print("basic consume loop")
|
print("basic consume loop")
|
||||||
tmq.subscribe(topic_list)
|
|
||||||
sub_list = tmq.subscription()
|
|
||||||
print("subscribed topics: ", sub_list)
|
|
||||||
c = 0
|
c = 0
|
||||||
l = 0
|
l = 0
|
||||||
for i in range(10):
|
for i in range(10):
|
||||||
|
@ -163,19 +162,22 @@ class TDTestCase:
|
||||||
res = tmq.poll(10)
|
res = tmq.poll(10)
|
||||||
print(f"loop {l}")
|
print(f"loop {l}")
|
||||||
l += 1
|
l += 1
|
||||||
if res:
|
if not res:
|
||||||
c += 1
|
|
||||||
topic = res.get_topic_name()
|
|
||||||
vg = res.get_vgroup_id()
|
|
||||||
db = res.get_db_name()
|
|
||||||
print(f"topic: {topic}\nvgroup id: {vg}\ndb: {db}")
|
|
||||||
for row in res:
|
|
||||||
print(row)
|
|
||||||
print("* committed")
|
|
||||||
tmq.commit(res)
|
|
||||||
else:
|
|
||||||
print(f"received empty message at loop {l} (committed {c})")
|
print(f"received empty message at loop {l} (committed {c})")
|
||||||
pass
|
continue
|
||||||
|
if res.error():
|
||||||
|
print(f"consumer error at loop {l} (committed {c}) {res.error()}")
|
||||||
|
continue
|
||||||
|
|
||||||
|
c += 1
|
||||||
|
topic = res.topic()
|
||||||
|
db = res.database()
|
||||||
|
print(f"topic: {topic}\ndb: {db}")
|
||||||
|
|
||||||
|
for row in res:
|
||||||
|
print(row.fetchall())
|
||||||
|
print("* committed")
|
||||||
|
tmq.commit(res)
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
tdSql.prepare()
|
tdSql.prepare()
|
||||||
|
@ -184,9 +186,11 @@ class TDTestCase:
|
||||||
self.drop_topic()
|
self.drop_topic()
|
||||||
self.user_privilege_check()
|
self.user_privilege_check()
|
||||||
self.subscribe_topic()
|
self.subscribe_topic()
|
||||||
|
|
||||||
def stop(self):
|
def stop(self):
|
||||||
tdSql.close()
|
tdSql.close()
|
||||||
tdLog.success("%s successfully executed" % __file__)
|
tdLog.success("%s successfully executed" % __file__)
|
||||||
|
|
||||||
|
|
||||||
tdCases.addWindows(__file__, TDTestCase())
|
tdCases.addWindows(__file__, TDTestCase())
|
||||||
tdCases.addLinux(__file__, TDTestCase())
|
tdCases.addLinux(__file__, TDTestCase())
|
Loading…
Reference in New Issue