Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
767d1cb2
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看板
提交
767d1cb2
编写于
1月 11, 2022
作者:
S
shenglian zhou
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
(query):check sversion/tversion when query rsp instead of fetch rsp
上级
9d41abcd
变更
4
显示空白变更内容
内联
并排
Showing
4 changed file
with
53 addition
and
70 deletion
+53
-70
src/client/src/tscServer.c
src/client/src/tscServer.c
+29
-33
src/inc/taosmsg.h
src/inc/taosmsg.h
+0
-2
src/query/src/queryMain.c
src/query/src/queryMain.c
+7
-31
src/vnode/src/vnodeRead.c
src/vnode/src/vnodeRead.c
+17
-4
未找到文件。
src/client/src/tscServer.c
浏览文件 @
767d1cb2
...
@@ -2704,14 +2704,29 @@ int tscProcessQueryRsp(SSqlObj *pSql) {
...
@@ -2704,14 +2704,29 @@ int tscProcessQueryRsp(SSqlObj *pSql) {
pRes
->
data
=
NULL
;
pRes
->
data
=
NULL
;
tscResetForNextRetrieve
(
pRes
);
tscResetForNextRetrieve
(
pRes
);
tscDebug
(
"0x%"
PRIx64
" query rsp received, qId:0x%"
PRIx64
,
pSql
->
self
,
pRes
->
qId
);
int32_t
sVersion
=
-
1
;
int32_t
tVersion
=
-
1
;
if
(
pQueryAttr
->
extend
==
1
)
{
STLV
*
tlv
=
(
STLV
*
)(
pRes
->
pRsp
+
sizeof
(
SQueryTableRsp
));
while
(
tlv
->
type
!=
TLV_TYPE_END_MARK
)
{
switch
(
ntohs
(
tlv
->
type
))
{
case
TLV_TYPE_META_VERSION
:
sVersion
=
ntohl
(
*
(
int32_t
*
)
tlv
->
value
);
tVersion
=
ntohl
(
*
(
int32_t
*
)(
tlv
->
value
+
sizeof
(
int32_t
)));
break
;
}
tlv
=
(
STLV
*
)
((
char
*
)
tlv
+
sizeof
(
STLV
)
+
ntohl
(
tlv
->
len
));
}
}
STableMetaInfo
*
tableMetaInfo
=
tscGetTableMetaInfoFromCmd
(
&
pSql
->
cmd
,
0
);
STableMetaInfo
*
tableMetaInfo
=
tscGetTableMetaInfoFromCmd
(
&
pSql
->
cmd
,
0
);
if
(
tableMetaInfo
->
pTableMeta
->
sversion
<
pQueryAttr
->
sVersion
||
if
(
tableMetaInfo
->
pTableMeta
->
sversion
<
sVersion
||
tableMetaInfo
->
pTableMeta
->
tversion
<
pQueryAttr
->
tVersion
)
{
tableMetaInfo
->
pTableMeta
->
tversion
<
tVersion
)
{
return
TSDB_CODE_TSC_INVALID_SCHEMA_VERSION
;
return
TSDB_CODE_TSC_INVALID_SCHEMA_VERSION
;
}
}
tscDebug
(
"0x%"
PRIx64
" query rsp received, qId:0x%"
PRIx64
,
pSql
->
self
,
pRes
->
qId
);
return
0
;
return
0
;
}
}
...
@@ -2802,16 +2817,16 @@ int tscProcessRetrieveRspFromNode(SSqlObj *pSql) {
...
@@ -2802,16 +2817,16 @@ int tscProcessRetrieveRspFromNode(SSqlObj *pSql) {
tscSetResRawPtr
(
pRes
,
pQueryInfo
,
pRes
->
dataConverted
);
tscSetResRawPtr
(
pRes
,
pQueryInfo
,
pRes
->
dataConverted
);
}
}
char
*
p
=
NULL
;
if
(
pSql
->
pSubscription
!=
NULL
)
{
if
(
pRetrieve
->
compressed
)
{
int32_t
numOfCols
=
pQueryInfo
->
fieldsInfo
.
numOfOutput
;
p
=
pRetrieve
->
data
+
ntohl
(
pRetrieve
->
compLen
)
+
pQueryInfo
->
fieldsInfo
.
numOfOutput
*
sizeof
(
int32_t
);
}
else
{
TAOS_FIELD
*
pField
=
tscFieldInfoGetField
(
&
pQueryInfo
->
fieldsInfo
,
numOfCols
-
1
);
p
=
pRetrieve
->
data
+
ntohl
(
pRetrieve
->
compLen
);
int16_t
offset
=
tscFieldInfoGetOffset
(
pQueryInfo
,
numOfCols
-
1
);
}
char
*
p
=
pRes
->
data
+
(
pField
->
bytes
+
offset
)
*
pRes
->
numOfRows
;
int32_t
numOfTables
=
htonl
(
*
(
int32_t
*
)
p
);
int32_t
numOfTables
=
htonl
(
*
(
int32_t
*
)
p
);
p
+=
sizeof
(
int32_t
);
p
+=
sizeof
(
int32_t
);
if
(
pSql
->
pSubscription
!=
NULL
)
{
for
(
int
i
=
0
;
i
<
numOfTables
;
i
++
)
{
for
(
int
i
=
0
;
i
<
numOfTables
;
i
++
)
{
int64_t
uid
=
htobe64
(
*
(
int64_t
*
)
p
);
int64_t
uid
=
htobe64
(
*
(
int64_t
*
)
p
);
p
+=
sizeof
(
int64_t
);
p
+=
sizeof
(
int64_t
);
...
@@ -2820,35 +2835,16 @@ int tscProcessRetrieveRspFromNode(SSqlObj *pSql) {
...
@@ -2820,35 +2835,16 @@ int tscProcessRetrieveRspFromNode(SSqlObj *pSql) {
p
+=
sizeof
(
TSKEY
);
p
+=
sizeof
(
TSKEY
);
tscUpdateSubscriptionProgress
(
pSql
->
pSubscription
,
uid
,
key
);
tscUpdateSubscriptionProgress
(
pSql
->
pSubscription
,
uid
,
key
);
}
}
}
else
{
p
+=
numOfTables
*
sizeof
(
STableIdInfo
);
}
}
pRes
->
row
=
0
;
pRes
->
row
=
0
;
tscDebug
(
"0x%"
PRIx64
" numOfRows:%d, offset:%"
PRId64
", complete:%d, qId:0x%"
PRIx64
,
pSql
->
self
,
pRes
->
numOfRows
,
pRes
->
offset
,
tscDebug
(
"0x%"
PRIx64
" numOfRows:%d, offset:%"
PRId64
", complete:%d, qId:0x%"
PRIx64
,
pSql
->
self
,
pRes
->
numOfRows
,
pRes
->
offset
,
pRes
->
completed
,
pRes
->
qId
);
pRes
->
completed
,
pRes
->
qId
);
if
(
pRetrieve
->
extend
==
1
)
{
STLV
*
tlv
=
(
STLV
*
)(
p
);
while
(
tlv
->
type
!=
TLV_TYPE_END_MARK
)
{
switch
(
ntohs
(
tlv
->
type
))
{
case
TLV_TYPE_META_VERSION
:
pRes
->
sVersion
=
ntohl
(
*
(
int32_t
*
)
tlv
->
value
);
pRes
->
tVersion
=
ntohl
(
*
(
int32_t
*
)(
tlv
->
value
+
sizeof
(
int32_t
)));
break
;
}
tlv
=
(
STLV
*
)
((
char
*
)
tlv
+
sizeof
(
STLV
)
+
ntohl
(
tlv
->
len
));
}
}
STableMeta
*
pTableMeta
=
pTableMetaInfo
->
pTableMeta
;
if
(
pTableMeta
->
sversion
<
pSql
->
res
.
sVersion
||
pTableMeta
->
tversion
<
pSql
->
res
.
tVersion
)
{
return
TSDB_CODE_TSC_INVALID_SCHEMA_VERSION
;
}
return
0
;
return
0
;
}
}
void
tscTableMetaCallBack
(
void
*
param
,
TAOS_RES
*
res
,
int
code
);
void
tscTableMetaCallBack
(
void
*
param
,
TAOS_RES
*
res
,
int
code
);
static
int32_t
getTableMetaFromMnode
(
SSqlObj
*
pSql
,
STableMetaInfo
*
pTableMetaInfo
,
bool
autocreate
)
{
static
int32_t
getTableMetaFromMnode
(
SSqlObj
*
pSql
,
STableMetaInfo
*
pTableMetaInfo
,
bool
autocreate
)
{
...
...
src/inc/taosmsg.h
浏览文件 @
767d1cb2
...
@@ -530,8 +530,6 @@ typedef struct {
...
@@ -530,8 +530,6 @@ typedef struct {
int8_t
extend
;
int8_t
extend
;
int32_t
code
;
int32_t
code
;
union
{
uint64_t
qhandle
;
uint64_t
qId
;};
// query handle
union
{
uint64_t
qhandle
;
uint64_t
qId
;};
// query handle
int32_t
sVersion
;
int32_t
tVersion
;
}
SQueryTableRsp
;
}
SQueryTableRsp
;
// todo: the show handle should be replaced with id
// todo: the show handle should be replaced with id
...
...
src/query/src/queryMain.c
浏览文件 @
767d1cb2
...
@@ -390,14 +390,11 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
...
@@ -390,14 +390,11 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
int32_t
s
=
GET_NUM_OF_RESULTS
(
pRuntimeEnv
);
int32_t
s
=
GET_NUM_OF_RESULTS
(
pRuntimeEnv
);
size_t
size
=
pQueryAttr
->
resultRowSize
*
s
;
size_t
size
=
pQueryAttr
->
resultRowSize
*
s
;
int32_t
contLenSubscriptions
=
(
int32_t
)(
size
+
sizeof
(
SRetrieveTableRsp
));
size
+=
sizeof
(
int32_t
);
size
+=
sizeof
(
int32_t
);
size
+=
sizeof
(
STableIdInfo
)
*
taosHashGetSize
(
pRuntimeEnv
->
pTableRetrieveTsMap
);
size
+=
sizeof
(
STableIdInfo
)
*
taosHashGetSize
(
pRuntimeEnv
->
pTableRetrieveTsMap
);
*
contLen
=
(
int32_t
)(
size
+
sizeof
(
SRetrieveTableRsp
));
*
contLen
=
(
int32_t
)(
size
+
sizeof
(
SRetrieveTableRsp
));
int32_t
contLenBeforeTLV
=
*
contLen
;
*
contLen
+=
(
sizeof
(
STLV
)
+
sizeof
(
int32_t
)
+
sizeof
(
int32_t
));
//tlv meta version
*
contLen
+=
sizeof
(
STLV
);
// tlv end mark
// current solution only avoid crash, but cannot return error code to client
// current solution only avoid crash, but cannot return error code to client
*
pRsp
=
(
SRetrieveTableRsp
*
)
rpcMallocCont
(
*
contLen
);
*
pRsp
=
(
SRetrieveTableRsp
*
)
rpcMallocCont
(
*
contLen
);
if
(
*
pRsp
==
NULL
)
{
if
(
*
pRsp
==
NULL
)
{
...
@@ -405,6 +402,7 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
...
@@ -405,6 +402,7 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
}
}
(
*
pRsp
)
->
numOfRows
=
htonl
((
int32_t
)
s
);
(
*
pRsp
)
->
numOfRows
=
htonl
((
int32_t
)
s
);
if
(
pQInfo
->
code
==
TSDB_CODE_SUCCESS
)
{
if
(
pQInfo
->
code
==
TSDB_CODE_SUCCESS
)
{
(
*
pRsp
)
->
offset
=
htobe64
(
pQInfo
->
runtimeEnv
.
currentOffset
);
(
*
pRsp
)
->
offset
=
htobe64
(
pQInfo
->
runtimeEnv
.
currentOffset
);
(
*
pRsp
)
->
useconds
=
htobe64
(
pQInfo
->
summary
.
elapsedTime
);
(
*
pRsp
)
->
useconds
=
htobe64
(
pQInfo
->
summary
.
elapsedTime
);
...
@@ -422,25 +420,16 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
...
@@ -422,25 +420,16 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
setQueryStatus
(
pRuntimeEnv
,
QUERY_OVER
);
setQueryStatus
(
pRuntimeEnv
,
QUERY_OVER
);
}
}
RESET_NUM_OF_RESULTS
(
&
(
pQInfo
->
runtimeEnv
));
if
((
*
pRsp
)
->
compressed
&&
compLen
!=
0
)
{
pQInfo
->
lastRetrieveTs
=
taosGetTimestampMs
();
int32_t
numOfCols
=
pQueryAttr
->
pExpr2
?
pQueryAttr
->
numOfExpr2
:
pQueryAttr
->
numOfOutput
;
int32_t
numOfCols
=
pQueryAttr
->
pExpr2
?
pQueryAttr
->
numOfExpr2
:
pQueryAttr
->
numOfOutput
;
int32_t
origSize
=
pQueryAttr
->
resultRowSize
*
s
;
int32_t
origSize
=
pQueryAttr
->
resultRowSize
*
s
;
int32_t
compSize
=
compLen
+
numOfCols
*
sizeof
(
int32_t
);
int32_t
compSize
=
compLen
+
numOfCols
*
sizeof
(
int32_t
);
if
((
*
pRsp
)
->
compressed
&&
compLen
!=
0
)
{
*
contLen
=
*
contLen
-
origSize
+
compSize
;
*
contLen
=
*
contLen
-
origSize
+
compSize
;
contLenSubscriptions
=
contLenSubscriptions
-
origSize
+
compSize
;
contLenBeforeTLV
=
contLenBeforeTLV
-
origSize
+
compSize
;
*
pRsp
=
(
SRetrieveTableRsp
*
)
rpcReallocCont
(
*
pRsp
,
*
contLen
);
*
pRsp
=
(
SRetrieveTableRsp
*
)
rpcReallocCont
(
*
pRsp
,
*
contLen
);
qDebug
(
"QInfo:0x%"
PRIx64
" compress col data, uncompressed size:%d, compressed size:%d, ratio:%.2f"
,
qDebug
(
"QInfo:0x%"
PRIx64
" compress col data, uncompressed size:%d, compressed size:%d, ratio:%.2f"
,
pQInfo
->
qId
,
origSize
,
compSize
,
(
float
)
origSize
/
(
float
)
compSize
);
pQInfo
->
qId
,
origSize
,
compSize
,
(
float
)
origSize
/
(
float
)
compSize
);
}
}
if
((
*
pRsp
)
->
compressed
)
{
(
*
pRsp
)
->
compLen
=
htonl
(
compLen
);
(
*
pRsp
)
->
compLen
=
htonl
(
compLen
);
}
else
{
(
*
pRsp
)
->
compLen
=
htonl
(
origSize
);
}
pQInfo
->
rspContext
=
NULL
;
pQInfo
->
rspContext
=
NULL
;
pQInfo
->
dataReady
=
QUERY_RESULT_NOT_READY
;
pQInfo
->
dataReady
=
QUERY_RESULT_NOT_READY
;
...
@@ -455,19 +444,6 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
...
@@ -455,19 +444,6 @@ int32_t qDumpRetrieveResult(qinfo_t qinfo, SRetrieveTableRsp **pRsp, int32_t *co
qDebug
(
"QInfo:0x%"
PRIx64
" has more results to retrieve"
,
pQInfo
->
qId
);
qDebug
(
"QInfo:0x%"
PRIx64
" has more results to retrieve"
,
pQInfo
->
qId
);
}
}
(
*
pRsp
)
->
extend
=
1
;
*
(
int32_t
*
)((
char
*
)(
*
pRsp
)
+
contLenSubscriptions
)
=
htonl
(
taosHashGetSize
(
pRuntimeEnv
->
pTableRetrieveTsMap
));
STLV
*
tlv
=
(
STLV
*
)((
char
*
)(
*
pRsp
)
+
contLenBeforeTLV
);
tlv
->
type
=
htons
(
TLV_TYPE_META_VERSION
);
tlv
->
len
=
htonl
(
sizeof
(
int32_t
)
+
sizeof
(
int32_t
));
int32_t
sVersion
=
htonl
(
pQueryAttr
->
tableGroupInfo
.
sVersion
);
int32_t
tVersion
=
htonl
(
pQueryAttr
->
tableGroupInfo
.
tVersion
);
memcpy
(
tlv
->
value
,
&
sVersion
,
sizeof
(
int32_t
));
memcpy
(
tlv
->
value
+
sizeof
(
int32_t
),
&
tVersion
,
sizeof
(
int32_t
));
STLV
*
tlvEnd
=
(
STLV
*
)((
char
*
)
tlv
+
sizeof
(
STLV
)
+
ntohl
(
tlv
->
len
));
tlvEnd
->
type
=
htons
(
TLV_TYPE_END_MARK
);
tlvEnd
->
len
=
0
;
// the memory should be freed if the code of pQInfo is not TSDB_CODE_SUCCESS
// the memory should be freed if the code of pQInfo is not TSDB_CODE_SUCCESS
if
(
pQInfo
->
code
!=
TSDB_CODE_SUCCESS
)
{
if
(
pQInfo
->
code
!=
TSDB_CODE_SUCCESS
)
{
rpcFreeCont
(
*
pRsp
);
rpcFreeCont
(
*
pRsp
);
...
...
src/vnode/src/vnodeRead.c
浏览文件 @
767d1cb2
...
@@ -243,14 +243,27 @@ static int32_t vnodeProcessQueryMsg(SVnodeObj *pVnode, SVReadMsg *pRead) {
...
@@ -243,14 +243,27 @@ static int32_t vnodeProcessQueryMsg(SVnodeObj *pVnode, SVReadMsg *pRead) {
int32_t
tagVersion
=
-
1
;
int32_t
tagVersion
=
-
1
;
code
=
qCreateQueryInfo
(
pVnode
->
tsdb
,
pVnode
->
vgId
,
pQueryTableMsg
,
&
pQInfo
,
qId
,
&
schemaVersion
,
&
tagVersion
);
code
=
qCreateQueryInfo
(
pVnode
->
tsdb
,
pVnode
->
vgId
,
pQueryTableMsg
,
&
pQInfo
,
qId
,
&
schemaVersion
,
&
tagVersion
);
SQueryTableRsp
*
pRsp
=
(
SQueryTableRsp
*
)
rpcMallocCont
(
sizeof
(
SQueryTableRsp
));
size_t
rspLen
=
sizeof
(
SQueryTableRsp
);
rspLen
+=
(
sizeof
(
STLV
)
+
sizeof
(
int32_t
)
+
sizeof
(
int32_t
));
rspLen
+=
sizeof
(
STLV
);
SQueryTableRsp
*
pRsp
=
(
SQueryTableRsp
*
)
rpcMallocCont
(
rspLen
);
pRsp
->
code
=
code
;
pRsp
->
code
=
code
;
pRsp
->
qId
=
0
;
pRsp
->
qId
=
0
;
pRsp
->
sVersion
=
schemaVersion
;
pRsp
->
extend
=
1
;
pRsp
->
tVersion
=
tagVersion
;
STLV
*
tlv
=
(
STLV
*
)((
char
*
)(
pRsp
)
+
sizeof
(
SQueryTableRsp
));
tlv
->
type
=
htons
(
TLV_TYPE_META_VERSION
);
tlv
->
len
=
htonl
(
sizeof
(
int32_t
)
+
sizeof
(
int32_t
));
int32_t
sVersion
=
htonl
(
schemaVersion
);
int32_t
tVersion
=
htonl
(
tagVersion
);
memcpy
(
tlv
->
value
,
&
sVersion
,
sizeof
(
int32_t
));
memcpy
(
tlv
->
value
+
sizeof
(
int32_t
),
&
tVersion
,
sizeof
(
int32_t
));
pRet
->
len
=
sizeof
(
SQueryTableRsp
);
STLV
*
tlvEnd
=
(
STLV
*
)((
char
*
)
tlv
+
sizeof
(
STLV
)
+
ntohl
(
tlv
->
len
));
tlvEnd
->
type
=
htons
(
TLV_TYPE_END_MARK
);
tlvEnd
->
len
=
0
;
pRet
->
len
=
rspLen
;
pRet
->
rsp
=
pRsp
;
pRet
->
rsp
=
pRsp
;
int32_t
vgId
=
pVnode
->
vgId
;
int32_t
vgId
=
pVnode
->
vgId
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录