Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
fc9c22b8
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
Fork
4786
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
fc9c22b8
编写于
5月 21, 2022
作者:
L
Liu Jicong
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(stream): array init error
上级
0a46e6ee
变更
1
隐藏空白更改
内联
并排
Showing
1 changed file
with
21 addition
and
23 deletion
+21
-23
source/libs/stream/src/tstreamUpdate.c
source/libs/stream/src/tstreamUpdate.c
+21
-23
未找到文件。
source/libs/stream/src/tstreamUpdate.c
浏览文件 @
fc9c22b8
...
...
@@ -17,21 +17,20 @@
#include "ttime.h"
#define DEFAULT_FALSE_POSITIVE 0.01
#define DEFAULT_BUCKET_SIZE 1024
#define ROWS_PER_MILLISECOND 1
#define MAX_NUM_SCALABLE_BF 120
#define MIN_NUM_SCALABLE_BF 10
#define DEFAULT_PREADD_BUCKET 1
#define MAX_INTERVAL MILLISECOND_PER_MINUTE
#define MIN_INTERVAL (MILLISECOND_PER_SECOND * 10)
#define DEFAULT_BUCKET_SIZE
1024
#define ROWS_PER_MILLISECOND
1
#define MAX_NUM_SCALABLE_BF
120
#define MIN_NUM_SCALABLE_BF
10
#define DEFAULT_PREADD_BUCKET
1
#define MAX_INTERVAL
MILLISECOND_PER_MINUTE
#define MIN_INTERVAL
(MILLISECOND_PER_SECOND * 10)
static
void
windowSBfAdd
(
SUpdateInfo
*
pInfo
,
uint64_t
count
)
{
if
(
pInfo
->
numSBFs
<
count
)
{
if
(
pInfo
->
numSBFs
<
count
)
{
count
=
pInfo
->
numSBFs
;
}
for
(
uint64_t
i
=
0
;
i
<
count
;
++
i
)
{
SScalableBf
*
tsSBF
=
tScalableBfInit
(
pInfo
->
interval
*
ROWS_PER_MILLISECOND
,
DEFAULT_FALSE_POSITIVE
);
SScalableBf
*
tsSBF
=
tScalableBfInit
(
pInfo
->
interval
*
ROWS_PER_MILLISECOND
,
DEFAULT_FALSE_POSITIVE
);
taosArrayPush
(
pInfo
->
pTsSBFs
,
&
tsSBF
);
}
}
...
...
@@ -76,7 +75,7 @@ static int64_t adjustWatermark(int64_t interval, int32_t watermark) {
return
watermark
;
}
SUpdateInfo
*
updateInfoInitP
(
SInterval
*
pInterval
,
int64_t
watermark
)
{
SUpdateInfo
*
updateInfoInitP
(
SInterval
*
pInterval
,
int64_t
watermark
)
{
return
updateInfoInit
(
pInterval
->
interval
,
pInterval
->
precision
,
watermark
);
}
...
...
@@ -93,7 +92,7 @@ SUpdateInfo *updateInfoInit(int64_t interval, int32_t precision, int64_t waterma
uint64_t
bfSize
=
(
uint64_t
)(
pInfo
->
watermark
/
pInfo
->
interval
);
pInfo
->
pTsSBFs
=
taosArrayInit
(
bfSize
,
sizeof
(
SScalableBf
));
pInfo
->
pTsSBFs
=
taosArrayInit
(
bfSize
,
sizeof
(
void
*
));
if
(
pInfo
->
pTsSBFs
==
NULL
)
{
updateInfoDestroy
(
pInfo
);
return
NULL
;
...
...
@@ -108,14 +107,14 @@ SUpdateInfo *updateInfoInit(int64_t interval, int32_t precision, int64_t waterma
}
TSKEY
dumy
=
0
;
for
(
uint64_t
i
=
0
;
i
<
DEFAULT_BUCKET_SIZE
;
++
i
)
{
for
(
uint64_t
i
=
0
;
i
<
DEFAULT_BUCKET_SIZE
;
++
i
)
{
taosArrayPush
(
pInfo
->
pTsBuckets
,
&
dumy
);
}
pInfo
->
numBuckets
=
DEFAULT_BUCKET_SIZE
;
return
pInfo
;
}
static
SScalableBf
*
getSBf
(
SUpdateInfo
*
pInfo
,
TSKEY
ts
)
{
static
SScalableBf
*
getSBf
(
SUpdateInfo
*
pInfo
,
TSKEY
ts
)
{
if
(
ts
<=
0
)
{
return
NULL
;
}
...
...
@@ -131,24 +130,23 @@ static SScalableBf* getSBf(SUpdateInfo *pInfo, TSKEY ts) {
}
SScalableBf
*
res
=
taosArrayGetP
(
pInfo
->
pTsSBFs
,
index
);
if
(
res
==
NULL
)
{
res
=
tScalableBfInit
(
pInfo
->
interval
*
ROWS_PER_MILLISECOND
,
DEFAULT_FALSE_POSITIVE
);
res
=
tScalableBfInit
(
pInfo
->
interval
*
ROWS_PER_MILLISECOND
,
DEFAULT_FALSE_POSITIVE
);
taosArrayPush
(
pInfo
->
pTsSBFs
,
&
res
);
}
return
res
;
}
bool
updateInfoIsUpdated
(
SUpdateInfo
*
pInfo
,
tb_uid_t
tableId
,
TSKEY
ts
)
{
int32_t
res
=
TSDB_CODE_FAILED
;
uint64_t
index
=
((
uint64_t
)
tableId
)
%
pInfo
->
numBuckets
;
SScalableBf
*
pSBf
=
getSBf
(
pInfo
,
ts
);
int32_t
res
=
TSDB_CODE_FAILED
;
uint64_t
index
=
((
uint64_t
)
tableId
)
%
pInfo
->
numBuckets
;
SScalableBf
*
pSBf
=
getSBf
(
pInfo
,
ts
);
// pSBf may be a null pointer
if
(
pSBf
)
{
res
=
tScalableBfPut
(
pSBf
,
&
ts
,
sizeof
(
TSKEY
));
}
TSKEY
maxTs
=
*
(
TSKEY
*
)
taosArrayGet
(
pInfo
->
pTsBuckets
,
index
);
if
(
maxTs
<
ts
)
{
if
(
maxTs
<
ts
)
{
taosArraySet
(
pInfo
->
pTsBuckets
,
index
,
&
ts
);
return
false
;
}
...
...
@@ -159,7 +157,7 @@ bool updateInfoIsUpdated(SUpdateInfo *pInfo, tb_uid_t tableId, TSKEY ts) {
return
false
;
}
//check from tsdb api
//
check from tsdb api
return
true
;
}
...
...
@@ -174,7 +172,7 @@ void updateInfoDestroy(SUpdateInfo *pInfo) {
SScalableBf
*
pSBF
=
taosArrayGetP
(
pInfo
->
pTsSBFs
,
i
);
tScalableBfDestroy
(
pSBF
);
}
taosArrayDestroy
(
pInfo
->
pTsSBFs
);
taosMemoryFree
(
pInfo
);
}
\ No newline at end of file
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录