Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
db1d75bb
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看板
提交
db1d75bb
编写于
4月 27, 2023
作者:
H
Haojun Liao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
refactor: do some internal refactor.
上级
a751f750
变更
13
隐藏空白更改
内联
并排
Showing
13 changed file
with
50 addition
and
56 deletion
+50
-56
include/common/tmsg.h
include/common/tmsg.h
+1
-1
source/common/src/tdatablock.c
source/common/src/tdatablock.c
+1
-1
source/common/src/tmsg.c
source/common/src/tmsg.c
+1
-1
source/dnode/vnode/inc/vnode.h
source/dnode/vnode/inc/vnode.h
+4
-7
source/dnode/vnode/src/sma/smaRollup.c
source/dnode/vnode/src/sma/smaRollup.c
+2
-2
source/dnode/vnode/src/sma/smaTimeRange.c
source/dnode/vnode/src/sma/smaTimeRange.c
+1
-1
source/dnode/vnode/src/tq/tqRead.c
source/dnode/vnode/src/tq/tqRead.c
+16
-18
source/dnode/vnode/src/tq/tqScan.c
source/dnode/vnode/src/tq/tqScan.c
+2
-2
source/dnode/vnode/src/tq/tqSink.c
source/dnode/vnode/src/tq/tqSink.c
+3
-3
source/dnode/vnode/src/vnd/vnodeSvr.c
source/dnode/vnode/src/vnd/vnodeSvr.c
+1
-1
source/libs/executor/src/dataInserter.c
source/libs/executor/src/dataInserter.c
+3
-3
source/libs/executor/src/scanoperator.c
source/libs/executor/src/scanoperator.c
+14
-15
source/libs/parser/src/parInsertUtil.c
source/libs/parser/src/parInsertUtil.c
+1
-1
未找到文件。
include/common/tmsg.h
浏览文件 @
db1d75bb
...
...
@@ -3423,7 +3423,7 @@ typedef struct {
int32_t
tEncodeSSubmitReq2
(
SEncoder
*
pCoder
,
const
SSubmitReq2
*
pReq
);
int32_t
tDecodeSSubmitReq2
(
SDecoder
*
pCoder
,
SSubmitReq2
*
pReq
);
void
tDestroySSubmitTbData
(
SSubmitTbData
*
pTbData
,
int32_t
flag
);
void
tDestroySSubmitReq
2
(
SSubmitReq2
*
pReq
,
int32_t
flag
);
void
tDestroySSubmitReq
(
SSubmitReq2
*
pReq
,
int32_t
flag
);
typedef
struct
{
int32_t
affectedRows
;
...
...
source/common/src/tdatablock.c
浏览文件 @
db1d75bb
...
...
@@ -2388,7 +2388,7 @@ _end:
if
(
terrno
!=
0
)
{
*
ppReq
=
NULL
;
if
(
pReq
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFreeClear
(
pReq
);
}
...
...
source/common/src/tmsg.c
浏览文件 @
db1d75bb
...
...
@@ -7436,7 +7436,7 @@ void tDestroySSubmitTbData(SSubmitTbData *pTbData, int32_t flag) {
}
}
void
tDestroySSubmitReq
2
(
SSubmitReq2
*
pReq
,
int32_t
flag
)
{
void
tDestroySSubmitReq
(
SSubmitReq2
*
pReq
,
int32_t
flag
)
{
if
(
pReq
->
aSubmitTbData
==
NULL
)
return
;
int32_t
nSubmitTbData
=
TARRAY_SIZE
(
pReq
->
aSubmitTbData
);
...
...
source/dnode/vnode/inc/vnode.h
浏览文件 @
db1d75bb
...
...
@@ -259,17 +259,14 @@ int32_t tqReaderAddTbUidList(STqReader *pReader, const SArray *pTableUidList);
int32_t
tqReaderRemoveTbUidList
(
STqReader
*
pReader
,
const
SArray
*
tbUidList
);
int32_t
tqSeekVer
(
STqReader
*
pReader
,
int64_t
ver
,
const
char
*
id
);
void
tqNextBlock
(
STqReader
*
pReader
,
SFetchRet
*
ret
);
int32_t
tqNextBlock
(
STqReader
*
pReader
,
SSDataBlock
*
pBlock
);
int32_t
extractSubmitMsgFromWal
(
SWalReader
*
pReader
,
SPackedData
*
pPackedData
);
int32_t
tqReaderSetSubmitMsg
(
STqReader
*
pReader
,
void
*
msgStr
,
int32_t
msgLen
,
int64_t
ver
);
// int32_t tqReaderSetDataMsg(STqReader *pReader, const SSubmitReq *pMsg, int64_t ver);
bool
tqNextDataBlock
(
STqReader
*
pReader
);
bool
tqNextDataBlockFilterOut2
(
STqReader
*
pReader
,
SHashObj
*
filterOutUids
);
int32_t
tqRetrieveDataBlock2
(
SSDataBlock
*
pBlock
,
STqReader
*
pReader
,
SSubmitTbData
**
pSubmitTbDataRet
);
bool
tqNextBlockImpl
(
STqReader
*
pReader
);
bool
tqNextDataBlockFilterOut
(
STqReader
*
pReader
,
SHashObj
*
filterOutUids
);
int32_t
tqRetrieveDataBlock
(
SSDataBlock
*
pBlock
,
STqReader
*
pReader
,
SSubmitTbData
**
pSubmitTbDataRet
);
int32_t
tqRetrieveTaosxBlock2
(
STqReader
*
pReader
,
SArray
*
blocks
,
SArray
*
schemas
,
SSubmitTbData
**
pSubmitTbDataRet
);
// int32_t tqRetrieveDataBlock(SSDataBlock *pBlock, STqReader *pReader);
// int32_t tqRetrieveTaosxBlock(STqReader *pReader, SArray *blocks, SArray *schemas);
int32_t
vnodeEnqueueStreamMsg
(
SVnode
*
pVnode
,
SRpcMsg
*
pMsg
);
...
...
source/dnode/vnode/src/sma/smaRollup.c
浏览文件 @
db1d75bb
...
...
@@ -684,7 +684,7 @@ static int32_t tdRSmaExecAndSubmitResult(SSma *pSma, qTaskInfo_t taskInfo, SRSma
}
if
(
pReq
&&
tdProcessSubmitReq
(
sinkTsdb
,
output
->
info
.
version
,
pReq
)
<
0
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
smaError
(
"vgId:%d, process submit req for rsma suid:%"
PRIu64
", uid:%"
PRIu64
" level %"
PRIi8
" failed since %s"
,
...
...
@@ -696,7 +696,7 @@ static int32_t tdRSmaExecAndSubmitResult(SSma *pSma, qTaskInfo_t taskInfo, SRSma
SMA_VID
(
pSma
),
suid
,
output
->
info
.
id
.
groupId
,
pItem
->
level
,
output
->
info
.
version
);
if
(
pReq
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
}
}
...
...
source/dnode/vnode/src/sma/smaTimeRange.c
浏览文件 @
db1d75bb
...
...
@@ -332,7 +332,7 @@ _end:
taosArrayDestroy
(
tagArray
);
taosArrayDestroy
(
pVals
);
if
(
pReq
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
}
...
...
source/dnode/vnode/src/tq/tqRead.c
浏览文件 @
db1d75bb
...
...
@@ -288,7 +288,7 @@ void tqCloseReader(STqReader* pReader) {
}
// free hash
taosHashCleanup
(
pReader
->
tbIdHash
);
tDestroySSubmitReq
2
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
tDestroySSubmitReq
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
taosMemoryFree
(
pReader
);
}
...
...
@@ -322,12 +322,11 @@ int32_t extractSubmitMsgFromWal(SWalReader* pReader, SPackedData* pPackedData) {
return
0
;
}
void
tqNextBlock
(
STqReader
*
pReader
,
SFetchRet
*
ret
)
{
int32_t
tqNextBlock
(
STqReader
*
pReader
,
SSDataBlock
*
pBlock
)
{
while
(
1
)
{
if
(
pReader
->
msg2
.
msgStr
==
NULL
)
{
if
(
walNextValidMsg
(
pReader
->
pWalReader
)
<
0
)
{
ret
->
fetchType
=
FETCH_TYPE__NONE
;
return
;
return
FETCH_TYPE__NONE
;
}
void
*
pBody
=
POINTER_SHIFT
(
pReader
->
pWalReader
->
pHead
->
head
.
body
,
sizeof
(
SSubmitReq2Msg
));
...
...
@@ -337,15 +336,14 @@ void tqNextBlock(STqReader* pReader, SFetchRet* ret) {
tqReaderSetSubmitMsg
(
pReader
,
pBody
,
bodyLen
,
ver
);
}
while
(
tqNext
DataBlock
(
pReader
))
{
memset
(
&
ret
->
data
,
0
,
sizeof
(
SSDataBlock
));
int32_t
code
=
tqRetrieveDataBlock
2
(
&
ret
->
data
,
pReader
,
NULL
);
if
(
code
!=
0
||
ret
->
data
.
info
.
rows
==
0
)
{
while
(
tqNext
BlockImpl
(
pReader
))
{
memset
(
pBlock
,
0
,
sizeof
(
SSDataBlock
));
int32_t
code
=
tqRetrieveDataBlock
(
pBlock
,
pReader
,
NULL
);
if
(
code
!=
TSDB_CODE_SUCCESS
||
pBlock
->
info
.
rows
==
0
)
{
continue
;
}
ret
->
fetchType
=
FETCH_TYPE__DATA
;
return
;
return
FETCH_TYPE__DATA
;
}
}
}
...
...
@@ -367,7 +365,7 @@ int32_t tqReaderSetSubmitMsg(STqReader* pReader, void* msgStr, int32_t msgLen, i
return
0
;
}
bool
tqNext
DataBlock
(
STqReader
*
pReader
)
{
bool
tqNext
BlockImpl
(
STqReader
*
pReader
)
{
if
(
pReader
->
msg2
.
msgStr
==
NULL
)
{
return
false
;
}
...
...
@@ -387,20 +385,20 @@ bool tqNextDataBlock(STqReader* pReader) {
tqDebug
(
"tq reader block found, ver:%"
PRId64
", uid:%"
PRId64
,
pReader
->
msg2
.
ver
,
pSubmitTbData
->
uid
);
return
true
;
}
else
{
tqDebug
(
"tq reader discard block, uid:%"
PRId64
", continue"
,
pSubmitTbData
->
uid
);
tqDebug
(
"tq reader discard
submit
block, uid:%"
PRId64
", continue"
,
pSubmitTbData
->
uid
);
}
pReader
->
nextBlk
++
;
}
tDestroySSubmitReq
2
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
tDestroySSubmitReq
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
pReader
->
nextBlk
=
0
;
pReader
->
msg2
.
msgStr
=
NULL
;
return
false
;
}
bool
tqNextDataBlockFilterOut
2
(
STqReader
*
pReader
,
SHashObj
*
filterOutUids
)
{
bool
tqNextDataBlockFilterOut
(
STqReader
*
pReader
,
SHashObj
*
filterOutUids
)
{
if
(
pReader
->
msg2
.
msgStr
==
NULL
)
return
false
;
int32_t
blockSz
=
taosArrayGetSize
(
pReader
->
submit
.
aSubmitTbData
);
...
...
@@ -415,7 +413,7 @@ bool tqNextDataBlockFilterOut2(STqReader* pReader, SHashObj* filterOutUids) {
pReader
->
nextBlk
++
;
}
tDestroySSubmitReq
2
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
tDestroySSubmitReq
(
&
pReader
->
submit
,
TSDB_MSG_FLG_DECODE
);
pReader
->
nextBlk
=
0
;
pReader
->
msg2
.
msgStr
=
NULL
;
...
...
@@ -451,7 +449,7 @@ int32_t tqMaskBlock(SSchemaWrapper* pDst, SSDataBlock* pBlock, const SSchemaWrap
return
0
;
}
int32_t
tqRetrieveDataBlock
2
(
SSDataBlock
*
pBlock
,
STqReader
*
pReader
,
SSubmitTbData
**
pSubmitTbDataRet
)
{
int32_t
tqRetrieveDataBlock
(
SSDataBlock
*
pBlock
,
STqReader
*
pReader
,
SSubmitTbData
**
pSubmitTbDataRet
)
{
tqDebug
(
"tq reader retrieve data block %p, index:%d"
,
pReader
->
msg2
.
msgStr
,
pReader
->
nextBlk
);
SSubmitTbData
*
pSubmitTbData
=
taosArrayGet
(
pReader
->
submit
.
aSubmitTbData
,
pReader
->
nextBlk
);
pReader
->
nextBlk
++
;
...
...
@@ -560,7 +558,7 @@ int32_t tqRetrieveDataBlock2(SSDataBlock* pBlock, STqReader* pReader, SSubmitTbD
int32_t
sourceIdx
=
0
;
while
(
targetIdx
<
colActual
)
{
if
(
sourceIdx
>=
numOfCols
){
tqError
(
"tqRetrieveDataBlock
2
sourceIdx:%d >= numOfCols:%d"
,
sourceIdx
,
numOfCols
);
tqError
(
"tqRetrieveDataBlock sourceIdx:%d >= numOfCols:%d"
,
sourceIdx
,
numOfCols
);
goto
FAIL
;
}
SColData
*
pCol
=
taosArrayGet
(
pCols
,
sourceIdx
);
...
...
@@ -568,7 +566,7 @@ int32_t tqRetrieveDataBlock2(SSDataBlock* pBlock, STqReader* pReader, SSubmitTbD
SColVal
colVal
;
if
(
pCol
->
nVal
!=
numOfRows
){
tqError
(
"tqRetrieveDataBlock
2
pCol->nVal:%d != numOfRows:%d"
,
pCol
->
nVal
,
numOfRows
);
tqError
(
"tqRetrieveDataBlock pCol->nVal:%d != numOfRows:%d"
,
pCol
->
nVal
,
numOfRows
);
goto
FAIL
;
}
...
...
source/dnode/vnode/src/tq/tqScan.c
浏览文件 @
db1d75bb
...
...
@@ -203,7 +203,7 @@ int32_t tqTaosxScanLog(STQ* pTq, STqHandle* pHandle, SPackedData submit, STaosxR
if
(
pExec
->
subType
==
TOPIC_SUB_TYPE__TABLE
)
{
STqReader
*
pReader
=
pExec
->
pTqReader
;
tqReaderSetSubmitMsg
(
pReader
,
submit
.
msgStr
,
submit
.
msgLen
,
submit
.
ver
);
while
(
tqNext
DataBlock
(
pReader
))
{
while
(
tqNext
BlockImpl
(
pReader
))
{
taosArrayClear
(
pBlocks
);
taosArrayClear
(
pSchemas
);
SSubmitTbData
*
pSubmitTbDataRet
=
NULL
;
...
...
@@ -262,7 +262,7 @@ int32_t tqTaosxScanLog(STQ* pTq, STqHandle* pHandle, SPackedData submit, STaosxR
}
else
if
(
pExec
->
subType
==
TOPIC_SUB_TYPE__DB
)
{
STqReader
*
pReader
=
pExec
->
pTqReader
;
tqReaderSetSubmitMsg
(
pReader
,
submit
.
msgStr
,
submit
.
msgLen
,
submit
.
ver
);
while
(
tqNextDataBlockFilterOut
2
(
pReader
,
pExec
->
execDb
.
pFilterOutTbUid
))
{
while
(
tqNextDataBlockFilterOut
(
pReader
,
pExec
->
execDb
.
pFilterOutTbUid
))
{
taosArrayClear
(
pBlocks
);
taosArrayClear
(
pSchemas
);
SSubmitTbData
*
pSubmitTbDataRet
=
NULL
;
...
...
source/dnode/vnode/src/tq/tqSink.c
浏览文件 @
db1d75bb
...
...
@@ -695,7 +695,7 @@ void tqSinkToTablePipeline2(SStreamTask* pTask, void* vnode, int64_t ver, void*
len
+=
sizeof
(
SSubmitReq2Msg
);
pBuf
=
rpcMallocCont
(
len
);
if
(
NULL
==
pBuf
)
{
tDestroySSubmitReq
2
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
goto
_end
;
}
((
SSubmitReq2Msg
*
)
pBuf
)
->
header
.
vgId
=
TD_VID
(
pVnode
);
...
...
@@ -707,11 +707,11 @@ void tqSinkToTablePipeline2(SStreamTask* pTask, void* vnode, int64_t ver, void*
tqError
(
"failed to encode submit req since %s"
,
terrstr
());
tEncoderClear
(
&
encoder
);
rpcFreeCont
(
pBuf
);
tDestroySSubmitReq
2
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
continue
;
}
tEncoderClear
(
&
encoder
);
tDestroySSubmitReq
2
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
&
submitReq
,
TSDB_MSG_FLG_ENCODE
);
SRpcMsg
msg
=
{
.
msgType
=
TDMT_VND_SUBMIT
,
...
...
source/dnode/vnode/src/vnd/vnodeSvr.c
浏览文件 @
db1d75bb
...
...
@@ -1388,7 +1388,7 @@ _exit:
// clear
taosArrayDestroy
(
newTbUids
);
tDestroySSubmitReq
2
(
pSubmitReq
,
0
==
pMsg
->
version
?
TSDB_MSG_FLG_CMPT
:
TSDB_MSG_FLG_DECODE
);
tDestroySSubmitReq
(
pSubmitReq
,
0
==
pMsg
->
version
?
TSDB_MSG_FLG_CMPT
:
TSDB_MSG_FLG_DECODE
);
tDestroySSubmitRsp2
(
pSubmitRsp
,
TSDB_MSG_FLG_ENCODE
);
if
(
code
)
terrno
=
code
;
...
...
source/libs/executor/src/dataInserter.c
浏览文件 @
db1d75bb
...
...
@@ -301,7 +301,7 @@ _end:
if
(
terrno
!=
0
)
{
*
ppReq
=
NULL
;
if
(
pReq
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
}
return
terrno
;
...
...
@@ -326,7 +326,7 @@ int32_t dataBlocksToSubmitReq(SDataInserterHandle* pInserter, void** pMsg, int32
code
=
buildSubmitReqFromBlock
(
pInserter
,
&
pReq
,
pDataBlock
,
pTSchema
,
uid
,
vgId
,
suid
);
if
(
code
)
{
if
(
pReq
)
{
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
}
...
...
@@ -335,7 +335,7 @@ int32_t dataBlocksToSubmitReq(SDataInserterHandle* pInserter, void** pMsg, int32
}
code
=
submitReqToMsg
(
vgId
,
pReq
,
pMsg
,
msgLen
);
tDestroySSubmitReq
2
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pReq
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pReq
);
return
code
;
...
...
source/libs/executor/src/scanoperator.c
浏览文件 @
db1d75bb
...
...
@@ -1646,10 +1646,10 @@ static SSDataBlock* doQueueScan(SOperatorInfo* pOperator) {
blockDataCleanup
(
pInfo
->
pRes
);
SDataBlockInfo
*
pBlockInfo
=
&
pInfo
->
pRes
->
info
;
while
(
tqNext
DataBlock
(
pInfo
->
tqReader
))
{
while
(
tqNext
BlockImpl
(
pInfo
->
tqReader
))
{
SSDataBlock
block
=
{
0
};
int32_t
code
=
tqRetrieveDataBlock
2
(
&
block
,
pInfo
->
tqReader
,
NULL
);
int32_t
code
=
tqRetrieveDataBlock
(
&
block
,
pInfo
->
tqReader
,
NULL
);
if
(
code
!=
TSDB_CODE_SUCCESS
||
block
.
info
.
rows
==
0
)
{
continue
;
}
...
...
@@ -1687,23 +1687,23 @@ static SSDataBlock* doQueueScan(SOperatorInfo* pOperator) {
if
(
pTaskInfo
->
streamInfo
.
currentOffset
.
type
==
TMQ_OFFSET__LOG
)
{
while
(
1
)
{
S
FetchRet
ret
=
{
0
};
tqNextBlock
(
pInfo
->
tqReader
,
&
ret
);
tqOffsetResetToLog
(
&
pTaskInfo
->
streamInfo
.
currentOffset
,
pInfo
->
tqReader
->
pWalReader
->
curVersion
-
1
);
// curVersion move to next, so currentOffset = curVersion - 1
if
(
ret
.
fetchT
ype
==
FETCH_TYPE__DATA
)
{
qDebug
(
"doQueueScan get data from log %"
PRId64
" rows, version:%"
PRId64
,
ret
.
data
.
info
.
rows
,
S
SDataBlock
block
=
{
0
};
int32_t
type
=
tqNextBlock
(
pInfo
->
tqReader
,
&
block
);
// curVersion move to next, so currentOffset = curVersion - 1
tqOffsetResetToLog
(
&
pTaskInfo
->
streamInfo
.
currentOffset
,
pInfo
->
tqReader
->
pWalReader
->
curVersion
-
1
);
if
(
t
ype
==
FETCH_TYPE__DATA
)
{
qDebug
(
"doQueueScan get data from log %"
PRId64
" rows, version:%"
PRId64
,
block
.
info
.
rows
,
pTaskInfo
->
streamInfo
.
currentOffset
.
version
);
blockDataCleanup
(
pInfo
->
pRes
);
setBlockIntoRes
(
pInfo
,
&
ret
.
data
,
true
);
setBlockIntoRes
(
pInfo
,
&
block
,
true
);
if
(
pInfo
->
pRes
->
info
.
rows
>
0
)
{
qDebug
(
"doQueueScan get data from log %"
PRId64
" rows, return, version:%"
PRId64
,
pInfo
->
pRes
->
info
.
rows
,
pTaskInfo
->
streamInfo
.
currentOffset
.
version
);
return
pInfo
->
pRes
;
}
}
else
if
(
ret
.
fetchT
ype
==
FETCH_TYPE__NONE
)
{
}
else
if
(
t
ype
==
FETCH_TYPE__NONE
)
{
qDebug
(
"doQueueScan get none from log, return, version:%"
PRId64
,
pTaskInfo
->
streamInfo
.
currentOffset
.
version
);
return
NULL
;
}
...
...
@@ -2072,11 +2072,10 @@ FETCH_NEXT_BLOCK:
blockDataCleanup
(
pInfo
->
pRes
);
while
(
tqNext
DataBlock
(
pInfo
->
tqReader
))
{
while
(
tqNext
BlockImpl
(
pInfo
->
tqReader
))
{
SSDataBlock
block
=
{
0
};
int32_t
code
=
tqRetrieveDataBlock2
(
&
block
,
pInfo
->
tqReader
,
NULL
);
int32_t
code
=
tqRetrieveDataBlock
(
&
block
,
pInfo
->
tqReader
,
NULL
);
if
(
code
!=
TSDB_CODE_SUCCESS
||
block
.
info
.
rows
==
0
)
{
continue
;
}
...
...
source/libs/parser/src/parInsertUtil.c
浏览文件 @
db1d75bb
...
...
@@ -324,7 +324,7 @@ void insDestroyVgroupDataCxt(SVgroupDataCxt* pVgCxt) {
return
;
}
tDestroySSubmitReq
2
(
pVgCxt
->
pData
,
TSDB_MSG_FLG_ENCODE
);
tDestroySSubmitReq
(
pVgCxt
->
pData
,
TSDB_MSG_FLG_ENCODE
);
taosMemoryFree
(
pVgCxt
->
pData
);
taosMemoryFree
(
pVgCxt
);
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录