Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
e1046015
T
TDengine
项目概览
taosdata
/
TDengine
大约 2 年 前同步成功
通知
1192
Star
22018
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看板
提交
e1046015
编写于
3月 28, 2023
作者:
dengyihao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
add backend
上级
47bd13a7
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
32 addition
and
16 deletion
+32
-16
source/libs/executor/src/executorimpl.c
source/libs/executor/src/executorimpl.c
+3
-3
source/libs/stream/src/streamStateRocksdb.c
source/libs/stream/src/streamStateRocksdb.c
+29
-13
未找到文件。
source/libs/executor/src/executorimpl.c
浏览文件 @
e1046015
...
...
@@ -2709,7 +2709,7 @@ int32_t getOutputBuf(SStreamState* pState, SWinKey* pKey, SResultRow** pResult,
}
int32_t
streamStateAddIfNotExist2
(
SStreamState
*
pState
,
const
SWinKey
*
key
,
void
**
pVal
,
int32_t
*
pVLen
)
{
qWarn
(
"streamStateAddIfNotExist"
);
//
qWarn("streamStateAddIfNotExist");
char
*
tVal
=
NULL
;
int32_t
size
=
0
;
int32_t
code
=
streamStateGet
(
pState
,
key
,
(
void
**
)
&
tVal
,
&
size
);
...
...
@@ -2803,8 +2803,8 @@ int32_t buildDataBlockFromGroupRes(SOperatorInfo* pOperator, SStreamState* pStat
pCtx
[
j
].
resultInfo
=
getResultEntryInfo
(
pRow
,
j
,
rowEntryOffset
);
SResultRowEntryInfo
*
pEnryInfo
=
pCtx
[
j
].
resultInfo
;
q
Warn
(
"initd:%d, complete:%d, null:%d, res:%d"
,
pEnryInfo
->
initialized
,
pEnryInfo
->
complete
,
pEnryInfo
->
isNullRes
,
pEnryInfo
->
numOfRes
);
q
Debug
(
"initd:%d, complete:%d, null:%d, res:%d"
,
pEnryInfo
->
initialized
,
pEnryInfo
->
complete
,
pEnryInfo
->
isNullRes
,
pEnryInfo
->
numOfRes
);
if
(
pCtx
[
j
].
fpSet
.
finalize
)
{
int32_t
code1
=
pCtx
[
j
].
fpSet
.
finalize
(
&
pCtx
[
j
],
pBlock
);
if
(
TAOS_FAILED
(
code1
))
{
...
...
source/libs/stream/src/streamStateRocksdb.c
浏览文件 @
e1046015
...
...
@@ -451,6 +451,7 @@ rocksdb_iterator_t* streamStateIterCreate(SStreamState* pState, const char* cfNa
char* val = rocksdb_get_cf(db, opts, pHandle, (const char*)buf, sizeof(*key), (size_t*)&len, &err); \
if (val == NULL) { \
qWarn("streamState str: %s failed to read from %s, err: not exist", toString, funcname); \
if (err != NULL) taosMemoryFree(err); \
code = -1; \
} else { \
if (pVal != NULL) *pVal = val; \
...
...
@@ -551,20 +552,35 @@ int32_t streamStateFillDel_rocksdb(SStreamState* pState, const SWinKey* key) {
int32_t
streamStateClear_rocksdb
(
SStreamState
*
pState
)
{
qDebug
(
"streamStateClear_rocksdb"
);
SWinKey
key
=
{.
ts
=
0
,
.
groupId
=
0
};
// batch clear later
streamStatePut_rocksdb
(
pState
,
&
key
,
NULL
,
0
);
while
(
1
)
{
SStreamStateCur
*
pCur
=
streamStateSeekKeyNext_rocksdb
(
pState
,
&
key
);
SWinKey
delKey
=
{
0
};
int32_t
code
=
streamStateGetKVByCur_rocksdb
(
pCur
,
&
delKey
,
NULL
,
0
);
streamStateFreeCur
(
pCur
);
if
(
code
==
0
)
{
streamStateDel_rocksdb
(
pState
,
&
delKey
);
}
else
{
break
;
}
SStateKey
sKey
=
{.
key
=
{.
ts
=
0
,
.
groupId
=
0
},
.
opNum
=
pState
->
number
};
SStateKey
eKey
=
{.
key
=
{.
ts
=
UINT64_MAX
,
.
groupId
=
INT64_MAX
},
.
opNum
=
pState
->
number
};
char
sKeyStr
[
128
]
=
{
0
};
char
eKeyStr
[
128
]
=
{
0
};
int
sLen
=
stateKeyEncode
(
&
sKey
,
sKeyStr
);
int
eLen
=
stateKeyEncode
(
&
sKey
,
eKeyStr
);
char
*
err
=
NULL
;
rocksdb_delete_range_cf
(
pState
->
pTdbState
->
rocksdb
,
pState
->
pTdbState
->
writeOpts
,
pState
->
pTdbState
->
pHandle
[
0
],
sKeyStr
,
sLen
,
eKeyStr
,
eLen
,
&
err
);
if
(
err
!=
NULL
)
{
qWarn
(
"failed to delete range cf(default)"
);
}
// batch clear later
// streamStatePut_rocksdb(pState, &key, NULL, 0);
// while (1) {
// SStreamStateCur* pCur = streamStateSeekKeyNext_rocksdb(pState, &key);
// SWinKey delKey = {0};
// int32_t code = streamStateGetKVByCur_rocksdb(pCur, &delKey, NULL, 0);
// streamStateFreeCur(pCur);
// if (code == 0) {
// streamStateDel_rocksdb(pState, &delKey);
// } else {
// break;
// }
// }
return
0
;
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录