Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
d7f7a7d6
T
TDengine
项目概览
taosdata
/
TDengine
大约 2 年 前同步成功
通知
1193
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看板
提交
d7f7a7d6
编写于
7月 17, 2023
作者:
H
Haojun Liao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix(stream): fix syntax error
上级
b2a905bd
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
1 addition
and
58 deletion
+1
-58
source/dnode/vnode/src/tq/tq.c
source/dnode/vnode/src/tq/tq.c
+0
-53
source/dnode/vnode/src/tq/tqUtil.c
source/dnode/vnode/src/tq/tqUtil.c
+0
-4
source/libs/stream/src/streamExec.c
source/libs/stream/src/streamExec.c
+1
-1
未找到文件。
source/dnode/vnode/src/tq/tq.c
浏览文件 @
d7f7a7d6
...
@@ -1330,59 +1330,6 @@ int32_t tqProcessTaskRecoverFinishRsp(STQ* pTq, SRpcMsg* pMsg) {
...
@@ -1330,59 +1330,6 @@ int32_t tqProcessTaskRecoverFinishRsp(STQ* pTq, SRpcMsg* pMsg) {
return
0
;
return
0
;
}
}
int32_t
extractDelDataBlock
(
const
void
*
pData
,
int32_t
len
,
int64_t
ver
,
SStreamRefDataBlock
**
pRefBlock
)
{
SDecoder
*
pCoder
=
&
(
SDecoder
){
0
};
SDeleteRes
*
pRes
=
&
(
SDeleteRes
){
0
};
(
*
pRefBlock
)
=
NULL
;
pRes
->
uidList
=
taosArrayInit
(
0
,
sizeof
(
tb_uid_t
));
if
(
pRes
->
uidList
==
NULL
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
tDecoderInit
(
pCoder
,
(
uint8_t
*
)
pData
,
len
);
tDecodeDeleteRes
(
pCoder
,
pRes
);
tDecoderClear
(
pCoder
);
int32_t
numOfTables
=
taosArrayGetSize
(
pRes
->
uidList
);
if
(
numOfTables
==
0
||
pRes
->
affectedRows
==
0
)
{
taosArrayDestroy
(
pRes
->
uidList
);
return
TSDB_CODE_SUCCESS
;
}
SSDataBlock
*
pDelBlock
=
createSpecialDataBlock
(
STREAM_DELETE_DATA
);
blockDataEnsureCapacity
(
pDelBlock
,
numOfTables
);
pDelBlock
->
info
.
rows
=
numOfTables
;
pDelBlock
->
info
.
version
=
ver
;
for
(
int32_t
i
=
0
;
i
<
numOfTables
;
i
++
)
{
// start key column
SColumnInfoData
*
pStartCol
=
taosArrayGet
(
pDelBlock
->
pDataBlock
,
START_TS_COLUMN_INDEX
);
colDataSetVal
(
pStartCol
,
i
,
(
const
char
*
)
&
pRes
->
skey
,
false
);
// end key column
SColumnInfoData
*
pEndCol
=
taosArrayGet
(
pDelBlock
->
pDataBlock
,
END_TS_COLUMN_INDEX
);
colDataSetVal
(
pEndCol
,
i
,
(
const
char
*
)
&
pRes
->
ekey
,
false
);
// uid column
SColumnInfoData
*
pUidCol
=
taosArrayGet
(
pDelBlock
->
pDataBlock
,
UID_COLUMN_INDEX
);
int64_t
*
pUid
=
taosArrayGet
(
pRes
->
uidList
,
i
);
colDataSetVal
(
pUidCol
,
i
,
(
const
char
*
)
pUid
,
false
);
colDataSetNULL
(
taosArrayGet
(
pDelBlock
->
pDataBlock
,
GROUPID_COLUMN_INDEX
),
i
);
colDataSetNULL
(
taosArrayGet
(
pDelBlock
->
pDataBlock
,
CALCULATE_START_TS_COLUMN_INDEX
),
i
);
colDataSetNULL
(
taosArrayGet
(
pDelBlock
->
pDataBlock
,
CALCULATE_END_TS_COLUMN_INDEX
),
i
);
}
taosArrayDestroy
(
pRes
->
uidList
);
*
pRefBlock
=
taosAllocateQitem
(
sizeof
(
SStreamRefDataBlock
),
DEF_QITEM
,
0
);
if
((
*
pRefBlock
)
==
NULL
)
{
return
TSDB_CODE_OUT_OF_MEMORY
;
}
(
*
pRefBlock
)
->
type
=
STREAM_INPUT__REF_DATA_BLOCK
;
(
*
pRefBlock
)
->
pBlock
=
pDelBlock
;
return
TSDB_CODE_SUCCESS
;
}
int32_t
tqProcessTaskRunReq
(
STQ
*
pTq
,
SRpcMsg
*
pMsg
)
{
int32_t
tqProcessTaskRunReq
(
STQ
*
pTq
,
SRpcMsg
*
pMsg
)
{
SStreamTaskRunReq
*
pReq
=
pMsg
->
pCont
;
SStreamTaskRunReq
*
pReq
=
pMsg
->
pCont
;
...
...
source/dnode/vnode/src/tq/tqUtil.c
浏览文件 @
d7f7a7d6
...
@@ -454,7 +454,3 @@ int32_t extractDelDataBlock(const void* pData, int32_t len, int64_t ver, SStream
...
@@ -454,7 +454,3 @@ int32_t extractDelDataBlock(const void* pData, int32_t len, int64_t ver, SStream
(
*
pRefBlock
)
->
pBlock
=
pDelBlock
;
(
*
pRefBlock
)
->
pBlock
=
pDelBlock
;
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
int32_t
tqCreateCheckpointBlock
(
SStreamCheckpoint
**
pCheckpointBlock
)
{
return
0
;
}
source/libs/stream/src/streamExec.c
浏览文件 @
d7f7a7d6
...
@@ -598,7 +598,7 @@ int32_t streamTaskReloadState(SStreamTask* pTask) {
...
@@ -598,7 +598,7 @@ int32_t streamTaskReloadState(SStreamTask* pTask) {
}
}
int32_t
streamAlignTransferState
(
SStreamTask
*
pTask
)
{
int32_t
streamAlignTransferState
(
SStreamTask
*
pTask
)
{
int32_t
numOfUpstream
=
taosArrayGetSize
(
pTask
->
pUpstream
Ep
InfoList
);
int32_t
numOfUpstream
=
taosArrayGetSize
(
pTask
->
pUpstreamInfoList
);
int32_t
old
=
atomic_val_compare_exchange_32
(
&
pTask
->
transferStateAlignCnt
,
0
,
numOfUpstream
);
int32_t
old
=
atomic_val_compare_exchange_32
(
&
pTask
->
transferStateAlignCnt
,
0
,
numOfUpstream
);
if
(
old
==
0
)
{
if
(
old
==
0
)
{
qDebug
(
"s-task:%s set the transfer state aligncnt %d"
,
pTask
->
id
.
idStr
,
numOfUpstream
);
qDebug
(
"s-task:%s set the transfer state aligncnt %d"
,
pTask
->
id
.
idStr
,
numOfUpstream
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录