Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
慢慢CG
TDengine
提交
3e64ef25
T
TDengine
项目概览
慢慢CG
/
TDengine
与 Fork 源项目一致
Fork自
taosdata / TDengine
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
3e64ef25
编写于
3月 12, 2021
作者:
S
Shuduo Sang
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[TD-3192] <feature>: support stb limit and offset. generate data refactor.
上级
7f2116fb
变更
1
显示空白变更内容
内联
并排
Showing
1 changed file
with
163 addition
and
145 deletion
+163
-145
src/kit/taosdemo/taosdemo.c
src/kit/taosdemo/taosdemo.c
+163
-145
未找到文件。
src/kit/taosdemo/taosdemo.c
浏览文件 @
3e64ef25
...
...
@@ -98,6 +98,8 @@ extern char configDir[];
#define MAX_DATABASE_COUNT 256
#define INPUT_BUF_LEN 256
#define DEFAULT_TIMESTAMP_STEP 10
typedef
enum
CREATE_SUB_TALBE_MOD_EN
{
PRE_CREATE_SUBTBL
,
AUTO_CREATE_SUBTBL
,
...
...
@@ -3307,7 +3309,7 @@ static bool getMetaFromInsertJsonFile(cJSON* root) {
if
(
timestampStep
&&
timestampStep
->
type
==
cJSON_Number
)
{
g_Dbs
.
db
[
i
].
superTbls
[
j
].
timeStampStep
=
timestampStep
->
valueint
;
}
else
if
(
!
timestampStep
)
{
g_Dbs
.
db
[
i
].
superTbls
[
j
].
timeStampStep
=
1000
;
g_Dbs
.
db
[
i
].
superTbls
[
j
].
timeStampStep
=
DEFAULT_TIMESTAMP_STEP
;
}
else
{
printf
(
"ERROR: failed to read json, timestamp_step not found
\n
"
);
goto
PARSE_OVER
;
...
...
@@ -4344,89 +4346,27 @@ static int execInsert(threadInfo *winfo, char *buffer, int k)
return
affectedRows
;
}
// sync insertion
/*
1 thread: 100 tables * 2000 rows/s
1 thread: 10 tables * 20000 rows/s
6 thread: 300 tables * 2000 rows/s
2 taosinsertdata , 1 thread: 10 tables * 20000 rows/s
*/
static
void
*
syncWrite
(
void
*
sarg
)
{
threadInfo
*
winfo
=
(
threadInfo
*
)
sarg
;
SSuperTable
*
superTblInfo
=
winfo
->
superTblInfo
;
static
int
generateDataBuffer
(
int32_t
threadID
,
threadInfo
*
pThreadInfo
,
char
*
buffer
,
int64_t
insertRows
,
int64_t
startFrom
,
int64_t
startTime
,
int
*
pSampleUsePos
)
{
SSuperTable
*
superTblInfo
=
pThreadInfo
->
superTblInfo
;
int
ncols_per_record
=
1
;
// count first col ts
int
samplePos
=
0
;
if
(
superTblInfo
)
{
if
(
0
!=
prepareSampleData
(
superTblInfo
))
return
NULL
;
if
(
superTblInfo
->
numberOfTblInOneSql
>
0
)
{
syncWriteForNumberOfTblInOneSql
(
winfo
,
superTblInfo
->
sampleDataBuf
);
tmfree
(
superTblInfo
->
sampleDataBuf
);
return
NULL
;
}
}
else
{
if
(
superTblInfo
==
NULL
)
{
int
datatypeSeq
=
0
;
while
(
g_args
.
datatype
[
datatypeSeq
])
{
datatypeSeq
++
;
ncols_per_record
++
;
}
}
char
*
buffer
=
calloc
(
superTblInfo
?
superTblInfo
->
maxSqlLen
:
g_args
.
max_sql_len
,
1
);
if
(
NULL
==
buffer
)
{
fprintf
(
stderr
,
"Failed to calloc %d Bytes, reason:%s
\n
"
,
superTblInfo
->
maxSqlLen
,
strerror
(
errno
));
tmfree
(
superTblInfo
->
sampleDataBuf
);
return
NULL
;
}
int64_t
lastPrintTime
=
taosGetTimestampMs
();
int64_t
startTs
=
taosGetTimestampUs
();
int64_t
endTs
;
int
insert_interval
=
superTblInfo
?
superTblInfo
->
insertInterval
:
g_args
.
insert_interval
;
uint64_t
st
=
0
;
uint64_t
et
=
0xffffffff
;
winfo
->
totalInsertRows
=
0
;
winfo
->
totalAffectedRows
=
0
;
int
sampleUsePos
;
if
(
superTblInfo
&&
superTblInfo
->
childTblLimit
)
{
// TODO
}
for
(
uint32_t
tID
=
winfo
->
start_table_id
;
tID
<=
winfo
->
end_table_id
;
tID
++
)
{
int64_t
start_time
=
winfo
->
start_time
;
int64_t
insertRows
=
(
superTblInfo
)
?
superTblInfo
->
insertRows
:
g_args
.
num_of_DPT
;
verbosePrint
(
"%s() LN%d insertRows=%"
PRId64
"
\n
"
,
__func__
,
__LINE__
,
insertRows
);
for
(
int64_t
i
=
0
;
i
<
insertRows
;)
{
int64_t
prepared
=
i
;
if
(
insert_interval
)
{
st
=
taosGetTimestampUs
();
}
sampleUsePos
=
samplePos
;
memset
(
buffer
,
0
,
superTblInfo
?
superTblInfo
->
maxSqlLen
:
g_args
.
max_sql_len
);
char
*
pstr
=
buffer
;
if
(
superTblInfo
)
{
if
(
AUTO_CREATE_SUBTBL
==
superTblInfo
->
autoCreateTable
)
{
char
*
tagsValBuf
=
NULL
;
if
(
0
==
superTblInfo
->
tagSource
)
{
...
...
@@ -4434,19 +4374,20 @@ static void* syncWrite(void *sarg) {
}
else
{
tagsValBuf
=
getTagValueFromTagSample
(
superTblInfo
,
tID
%
superTblInfo
->
tagSampleCount
);
t
hread
ID
%
superTblInfo
->
tagSampleCount
);
}
if
(
NULL
==
tagsValBuf
)
{
goto
free_and_statistics_2
;
fprintf
(
stderr
,
"tag buf failed to allocate memory
\n
"
);
return
-
1
;
}
pstr
+=
snprintf
(
pstr
,
superTblInfo
->
maxSqlLen
,
"insert into %s.%s%d using %s.%s tags %s values"
,
wi
nfo
->
db_name
,
pThreadI
nfo
->
db_name
,
superTblInfo
->
childTblPrefix
,
tID
,
wi
nfo
->
db_name
,
t
hread
ID
,
pThreadI
nfo
->
db_name
,
superTblInfo
->
sTblName
,
tagsValBuf
);
tmfree
(
tagsValBuf
);
...
...
@@ -4454,24 +4395,23 @@ static void* syncWrite(void *sarg) {
pstr
+=
snprintf
(
pstr
,
superTblInfo
->
maxSqlLen
,
"insert into %s.%s values"
,
wi
nfo
->
db_name
,
superTblInfo
->
childTblName
+
tID
*
TSDB_TABLE_NAME_LEN
);
pThreadI
nfo
->
db_name
,
superTblInfo
->
childTblName
+
t
hread
ID
*
TSDB_TABLE_NAME_LEN
);
}
else
{
pstr
+=
snprintf
(
pstr
,
(
superTblInfo
?
superTblInfo
->
maxSqlLen
:
g_args
.
max_sql_len
),
"insert into %s.%s%d values"
,
wi
nfo
->
db_name
,
pThreadI
nfo
->
db_name
,
superTblInfo
?
superTblInfo
->
childTblPrefix
:
g_args
.
tb_prefix
,
tID
);
t
hread
ID
);
}
}
else
{
pstr
+=
snprintf
(
pstr
,
(
superTblInfo
?
superTblInfo
->
maxSqlLen
:
g_args
.
max_sql_len
),
"insert into %s.%s%d values"
,
wi
nfo
->
db_name
,
pThreadI
nfo
->
db_name
,
superTblInfo
?
superTblInfo
->
childTblPrefix
:
g_args
.
tb_prefix
,
tID
);
t
hread
ID
);
}
int
k
;
...
...
@@ -4479,7 +4419,6 @@ static void* syncWrite(void *sarg) {
verbosePrint
(
"%s() LN%d num_of_RPR=%d
\n
"
,
__func__
,
__LINE__
,
g_args
.
num_of_RPR
);
for
(
k
=
0
;
k
<
g_args
.
num_of_RPR
;)
{
if
(
superTblInfo
)
{
int
retLen
=
0
;
...
...
@@ -4487,14 +4426,14 @@ static void* syncWrite(void *sarg) {
retLen
=
getRowDataFromSample
(
pstr
+
len
,
superTblInfo
->
maxSqlLen
-
len
,
start
_time
+
superTblInfo
->
timeStampStep
*
i
,
start
Time
+
superTblInfo
->
timeStampStep
*
startFrom
,
superTblInfo
,
&
s
ampleUsePos
);
pS
ampleUsePos
);
}
else
if
(
0
==
strncasecmp
(
superTblInfo
->
dataSource
,
"rand"
,
strlen
(
"rand"
)))
{
int
rand_num
=
rand_tinyint
()
%
100
;
if
(
0
!=
superTblInfo
->
disorderRatio
&&
rand_num
<
superTblInfo
->
disorderRatio
)
{
int64_t
d
=
start
_t
ime
-
rand
()
%
superTblInfo
->
disorderRange
;
int64_t
d
=
start
T
ime
-
rand
()
%
superTblInfo
->
disorderRange
;
retLen
=
generateRowData
(
pstr
+
len
,
superTblInfo
->
maxSqlLen
-
len
,
...
...
@@ -4505,12 +4444,12 @@ static void* syncWrite(void *sarg) {
retLen
=
generateRowData
(
pstr
+
len
,
superTblInfo
->
maxSqlLen
-
len
,
start
_time
+
superTblInfo
->
timeStampStep
*
i
,
start
Time
+
superTblInfo
->
timeStampStep
*
startFrom
,
superTblInfo
);
}
if
(
retLen
<
0
)
{
goto
free_and_statistics_2
;
return
-
1
;
}
len
+=
retLen
;
...
...
@@ -4524,12 +4463,14 @@ static void* syncWrite(void *sarg) {
if
((
g_args
.
disorderRatio
!=
0
)
&&
(
rand_num
<
g_args
.
disorderRange
))
{
int64_t
d
=
start_t
ime
-
rand
()
%
1000000
+
rand_num
;
int64_t
d
=
startT
ime
-
rand
()
%
1000000
+
rand_num
;
len
=
generateData
(
data
,
data_type
,
ncols_per_record
,
d
,
lenOfBinary
);
}
else
{
len
=
generateData
(
data
,
data_type
,
ncols_per_record
,
start_time
+=
1000
,
lenOfBinary
);
ncols_per_record
,
startTime
+
DEFAULT_TIMESTAMP_STEP
*
startFrom
,
lenOfBinary
);
}
//assert(len + pstr - buffer < BUFFER_SIZE);
...
...
@@ -4542,19 +4483,96 @@ static void* syncWrite(void *sarg) {
verbosePrint
(
"%s() LN%d len=%d k=%d
\n
buffer=%s
\n
"
,
__func__
,
__LINE__
,
len
,
k
,
buffer
);
prepared
++
;
k
++
;
i
++
;
startFrom
++
;
if
(
prepared
>=
insertRows
)
debugPrint
(
"%s() LN%d k=%d startFrom=%ld insertRows=%ld
\n
"
,
__func__
,
__LINE__
,
k
,
startFrom
,
insertRows
);
if
(
startFrom
>=
insertRows
)
break
;
}
int
affectedRows
=
execInsert
(
winfo
,
buffer
,
k
);
return
k
;
}
// sync insertion
/*
1 thread: 100 tables * 2000 rows/s
1 thread: 10 tables * 20000 rows/s
6 thread: 300 tables * 2000 rows/s
2 taosinsertdata , 1 thread: 10 tables * 20000 rows/s
*/
static
void
*
syncWrite
(
void
*
sarg
)
{
threadInfo
*
winfo
=
(
threadInfo
*
)
sarg
;
SSuperTable
*
superTblInfo
=
winfo
->
superTblInfo
;
if
(
superTblInfo
)
{
if
(
0
!=
prepareSampleData
(
superTblInfo
))
return
NULL
;
if
(
superTblInfo
->
numberOfTblInOneSql
>
0
)
{
syncWriteForNumberOfTblInOneSql
(
winfo
,
superTblInfo
->
sampleDataBuf
);
tmfree
(
superTblInfo
->
sampleDataBuf
);
return
NULL
;
}
}
int
samplePos
=
0
;
char
*
buffer
=
calloc
(
superTblInfo
?
superTblInfo
->
maxSqlLen
:
g_args
.
max_sql_len
,
1
);
if
(
NULL
==
buffer
)
{
fprintf
(
stderr
,
"Failed to calloc %d Bytes, reason:%s
\n
"
,
superTblInfo
->
maxSqlLen
,
strerror
(
errno
));
tmfree
(
superTblInfo
->
sampleDataBuf
);
return
NULL
;
}
int64_t
lastPrintTime
=
taosGetTimestampMs
();
int64_t
startTs
=
taosGetTimestampUs
();
int64_t
endTs
;
int
insert_interval
=
superTblInfo
?
superTblInfo
->
insertInterval
:
g_args
.
insert_interval
;
uint64_t
st
=
0
;
uint64_t
et
=
0xffffffff
;
winfo
->
totalInsertRows
=
0
;
winfo
->
totalAffectedRows
=
0
;
int
sampleUsePos
;
if
(
superTblInfo
&&
superTblInfo
->
childTblLimit
)
{
// TODO
}
for
(
uint32_t
tID
=
winfo
->
start_table_id
;
tID
<=
winfo
->
end_table_id
;
tID
++
)
{
int64_t
start_time
=
winfo
->
start_time
;
int64_t
insertRows
=
(
superTblInfo
)
?
superTblInfo
->
insertRows
:
g_args
.
num_of_DPT
;
verbosePrint
(
"%s() LN%d insertRows=%"
PRId64
"
\n
"
,
__func__
,
__LINE__
,
insertRows
);
for
(
int64_t
i
=
0
;
i
<
insertRows
;)
{
if
(
insert_interval
)
{
st
=
taosGetTimestampUs
();
}
sampleUsePos
=
samplePos
;
int
generated
=
generateDataBuffer
(
tID
,
winfo
,
buffer
,
insertRows
,
i
,
start_time
,
&
sampleUsePos
);
if
(
generated
>
0
)
i
+=
generated
;
else
goto
free_and_statistics_2
;
int
affectedRows
=
execInsert
(
winfo
,
buffer
,
generated
);
if
(
affectedRows
<
0
)
goto
free_and_statistics_2
;
winfo
->
totalInsertRows
+=
k
;
winfo
->
totalInsertRows
+=
generated
;
winfo
->
totalAffectedRows
+=
affectedRows
;
endTs
=
taosGetTimestampUs
();
...
...
@@ -4573,7 +4591,7 @@ static void* syncWrite(void *sarg) {
lastPrintTime
=
currentPrintTime
;
}
if
(
prepared
>=
insertRows
)
if
(
i
>=
insertRows
)
break
;
if
(
insert_interval
)
{
...
...
@@ -4581,7 +4599,7 @@ static void* syncWrite(void *sarg) {
if
(
insert_interval
>
((
et
-
st
)
/
1000
)
)
{
int
sleep_time
=
insert_interval
-
(
et
-
st
)
/
1000
;
printf
(
"sleep: %d ms for insert interval
\n
"
,
sleep_time
);
verbosePrint
(
"%s() LN%d sleep: %d ms for insert interval
\n
"
,
__func__
,
__LINE__
,
sleep_time
);
taosMsleep
(
sleep_time
);
// ms
}
}
...
...
@@ -5675,7 +5693,7 @@ void setParaFromArg(){
tstrncpy
(
g_Dbs
.
db
[
0
].
superTbls
[
0
].
insertMode
,
"taosc"
,
MAX_TB_NAME_SIZE
);
tstrncpy
(
g_Dbs
.
db
[
0
].
superTbls
[
0
].
startTimestamp
,
"2017-07-14 10:40:00.000"
,
MAX_TB_NAME_SIZE
);
g_Dbs
.
db
[
0
].
superTbls
[
0
].
timeStampStep
=
10
;
g_Dbs
.
db
[
0
].
superTbls
[
0
].
timeStampStep
=
DEFAULT_TIMESTAMP_STEP
;
g_Dbs
.
db
[
0
].
superTbls
[
0
].
insertRows
=
g_args
.
num_of_DPT
;
g_Dbs
.
db
[
0
].
superTbls
[
0
].
maxSqlLen
=
TSDB_PAYLOAD_SIZE
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录