Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
慢慢CG
TDengine
提交
bc8cffd2
T
TDengine
项目概览
慢慢CG
/
TDengine
与 Fork 源项目一致
Fork自
taosdata / TDengine
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
bc8cffd2
编写于
11月 04, 2020
作者:
H
Haojun Liao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[TD-1373]<fix>: fix bugs in super table join.
上级
f3425fc9
变更
8
隐藏空白更改
内联
并排
Showing
8 changed file
with
196 addition
and
134 deletion
+196
-134
src/client/inc/tscSubquery.h
src/client/inc/tscSubquery.h
+1
-1
src/client/inc/tscUtil.h
src/client/inc/tscUtil.h
+1
-0
src/client/inc/tsclient.h
src/client/inc/tsclient.h
+3
-4
src/client/src/tscAsync.c
src/client/src/tscAsync.c
+2
-2
src/client/src/tscSQLParser.c
src/client/src/tscSQLParser.c
+35
-15
src/client/src/tscServer.c
src/client/src/tscServer.c
+3
-2
src/client/src/tscSubquery.c
src/client/src/tscSubquery.c
+139
-103
src/client/src/tscUtil.c
src/client/src/tscUtil.c
+12
-7
未找到文件。
src/client/inc/tscSubquery.h
浏览文件 @
bc8cffd2
...
...
@@ -23,7 +23,7 @@ extern "C" {
#include "tscUtil.h"
#include "tsclient.h"
void
tscFetchDatablockF
rom
Subquery
(
SSqlObj
*
pSql
);
void
tscFetchDatablockF
or
Subquery
(
SSqlObj
*
pSql
);
void
tscSetupOutputColumnIndex
(
SSqlObj
*
pSql
);
void
tscJoinQueryCallback
(
void
*
param
,
TAOS_RES
*
tres
,
int
code
);
...
...
src/client/inc/tscUtil.h
浏览文件 @
bc8cffd2
...
...
@@ -228,6 +228,7 @@ void tscClearSubqueryInfo(SSqlCmd* pCmd);
void
tscFreeVgroupTableInfo
(
SArray
*
pVgroupTables
);
SArray
*
tscVgroupTableInfoClone
(
SArray
*
pVgroupTables
);
void
tscRemoveVgroupTableGroup
(
SArray
*
pVgroupTable
,
int32_t
index
);
void
tscVgroupTableCopy
(
SVgroupTableInfo
*
info
,
SVgroupTableInfo
*
pInfo
);
int
tscGetSTableVgroupInfo
(
SSqlObj
*
pSql
,
int32_t
clauseIndex
);
int
tscGetTableMeta
(
SSqlObj
*
pSql
,
STableMetaInfo
*
pTableMetaInfo
);
...
...
src/client/inc/tsclient.h
浏览文件 @
bc8cffd2
...
...
@@ -339,9 +339,9 @@ typedef struct STscObj {
}
STscObj
;
typedef
struct
SSubqueryState
{
int32_t
numOfRemain
;
// the number of remain unfinished subquery
int32_t
numOfSub
;
// the number of total sub-queries
uint64_t
numOfRetrievedRows
;
// total number of points in this query
int32_t
numOfRemain
;
// the number of remain unfinished subquery
int32_t
numOfSub
;
// the number of total sub-queries
uint64_t
numOfRetrievedRows
;
// total number of points in this query
}
SSubqueryState
;
typedef
struct
SSqlObj
{
...
...
@@ -515,7 +515,6 @@ extern SRpcCorEpSet tscMgmtEpSet;
extern
int
(
*
tscBuildMsg
[
TSDB_SQL_MAX
])(
SSqlObj
*
pSql
,
SSqlInfo
*
pInfo
);
int32_t
tscCompareTidTags
(
const
void
*
p1
,
const
void
*
p2
);
void
tscBuildVgroupTableInfo
(
SSqlObj
*
pSql
,
STableMetaInfo
*
pTableMetaInfo
,
SArray
*
tables
);
#ifdef __cplusplus
...
...
src/client/src/tscAsync.c
浏览文件 @
bc8cffd2
...
...
@@ -176,7 +176,7 @@ static void tscProcessAsyncRetrieveImpl(void *param, TAOS_RES *tres, int numOfRo
}
if
(
pCmd
->
command
==
TSDB_SQL_TABLE_JOIN_RETRIEVE
)
{
tscFetchDatablockF
rom
Subquery
(
pSql
);
tscFetchDatablockF
or
Subquery
(
pSql
);
}
else
{
tscProcessSql
(
pSql
);
}
...
...
@@ -226,7 +226,7 @@ void taos_fetch_rows_a(TAOS_RES *taosa, __async_cb_func_t fp, void *param) {
// handle the sub queries of join query
if
(
pCmd
->
command
==
TSDB_SQL_TABLE_JOIN_RETRIEVE
)
{
tscFetchDatablockF
rom
Subquery
(
pSql
);
tscFetchDatablockF
or
Subquery
(
pSql
);
}
else
if
(
pRes
->
completed
)
{
if
(
pCmd
->
command
==
TSDB_SQL_FETCH
||
(
pCmd
->
command
>=
TSDB_SQL_SERV_STATUS
&&
pCmd
->
command
<=
TSDB_SQL_CURRENT_USER
))
{
if
(
hasMoreVnodesToTry
(
pSql
))
{
// sequentially retrieve data from remain vnodes.
...
...
src/client/src/tscSQLParser.c
浏览文件 @
bc8cffd2
...
...
@@ -1632,7 +1632,7 @@ int32_t addProjectionExprAndResultField(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, t
}
static
int32_t
setExprInfoForFunctions
(
SSqlCmd
*
pCmd
,
SQueryInfo
*
pQueryInfo
,
SSchema
*
pSchema
,
SConvertFunc
cvtFunc
,
tSQLExprItem
*
item
,
int32_t
resColIdx
,
SColumnIndex
*
pColIndex
,
bool
finalResult
)
{
const
char
*
name
,
int32_t
resColIdx
,
SColumnIndex
*
pColIndex
,
bool
finalResult
)
{
const
char
*
msg1
=
"not support column types"
;
int16_t
type
=
0
;
...
...
@@ -1640,7 +1640,7 @@ static int32_t setExprInfoForFunctions(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, SS
int32_t
functionID
=
cvtFunc
.
execFuncId
;
if
(
functionID
==
TSDB_FUNC_SPREAD
)
{
int32_t
t1
=
pSchema
[
pColIndex
->
columnIndex
].
type
;
int32_t
t1
=
pSchema
->
type
;
if
(
t1
==
TSDB_DATA_TYPE_BINARY
||
t1
==
TSDB_DATA_TYPE_NCHAR
||
t1
==
TSDB_DATA_TYPE_BOOL
)
{
invalidSqlErrMsg
(
tscGetErrorMsgPayload
(
pCmd
),
msg1
);
return
-
1
;
...
...
@@ -1649,17 +1649,12 @@ static int32_t setExprInfoForFunctions(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, SS
bytes
=
tDataTypeDesc
[
type
].
nSize
;
}
}
else
{
type
=
pSchema
[
pColIndex
->
columnIndex
].
type
;
bytes
=
pSchema
[
pColIndex
->
columnIndex
].
bytes
;
type
=
pSchema
->
type
;
bytes
=
pSchema
->
bytes
;
}
SSqlExpr
*
pExpr
=
tscSqlExprAppend
(
pQueryInfo
,
functionID
,
pColIndex
,
type
,
bytes
,
bytes
,
false
);
if
(
item
->
aliasName
!=
NULL
)
{
tstrncpy
(
pExpr
->
aliasName
,
item
->
aliasName
,
tListLen
(
pExpr
->
aliasName
));
}
else
{
int32_t
len
=
MIN
(
tListLen
(
pExpr
->
aliasName
),
item
->
pNode
->
token
.
n
+
1
);
tstrncpy
(
pExpr
->
aliasName
,
item
->
pNode
->
token
.
z
,
len
);
}
tstrncpy
(
pExpr
->
aliasName
,
name
,
tListLen
(
pExpr
->
aliasName
));
if
(
cvtFunc
.
originFuncId
==
TSDB_FUNC_LAST_ROW
&&
cvtFunc
.
originFuncId
!=
functionID
)
{
pExpr
->
colInfo
.
flag
|=
TSDB_COL_NULL
;
...
...
@@ -1687,6 +1682,18 @@ static int32_t setExprInfoForFunctions(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, SS
return
TSDB_CODE_SUCCESS
;
}
void
setResultColName
(
char
*
name
,
tSQLExprItem
*
pItem
,
int32_t
functionId
,
SStrToken
*
pToken
)
{
if
(
pItem
->
aliasName
!=
NULL
)
{
tstrncpy
(
name
,
pItem
->
aliasName
,
TSDB_COL_NAME_LEN
);
}
else
{
char
uname
[
TSDB_COL_NAME_LEN
]
=
{
0
};
int32_t
len
=
MIN
(
pToken
->
n
+
1
,
TSDB_COL_NAME_LEN
);
tstrncpy
(
uname
,
pToken
->
z
,
len
);
snprintf
(
name
,
TSDB_COL_NAME_LEN
-
1
,
"%s(%s)"
,
aAggs
[
functionId
].
aName
,
uname
);
}
}
int32_t
addExprAndResultField
(
SSqlCmd
*
pCmd
,
SQueryInfo
*
pQueryInfo
,
int32_t
colIndex
,
tSQLExprItem
*
pItem
,
bool
finalResult
)
{
STableMetaInfo
*
pTableMetaInfo
=
NULL
;
int32_t
optr
=
pItem
->
pNode
->
nSQLOptr
;
...
...
@@ -1954,9 +1961,13 @@ int32_t addExprAndResultField(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, int32_t col
pTableMetaInfo
=
tscGetMetaInfo
(
pQueryInfo
,
index
.
tableIndex
);
SSchema
*
pSchema
=
tscGetTableSchema
(
pTableMetaInfo
->
pTableMeta
);
char
name
[
TSDB_COL_NAME_LEN
]
=
{
0
};
for
(
int32_t
j
=
0
;
j
<
tscGetNumOfColumns
(
pTableMetaInfo
->
pTableMeta
);
++
j
)
{
index
.
columnIndex
=
j
;
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
pSchema
,
cvtFunc
,
pItem
,
colIndex
++
,
&
index
,
finalResult
)
!=
0
)
{
SStrToken
t
=
{.
z
=
pSchema
[
j
].
name
,
.
n
=
strnlen
(
pSchema
[
j
].
name
,
TSDB_COL_NAME_LEN
)};
setResultColName
(
name
,
pItem
,
cvtFunc
.
originFuncId
,
&
t
);
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
&
pSchema
[
j
],
cvtFunc
,
name
,
colIndex
++
,
&
index
,
finalResult
)
!=
0
)
{
return
TSDB_CODE_TSC_INVALID_SQL
;
}
}
...
...
@@ -1967,14 +1978,18 @@ int32_t addExprAndResultField(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, int32_t col
}
pTableMetaInfo
=
tscGetMetaInfo
(
pQueryInfo
,
index
.
tableIndex
);
SSchema
*
pSchema
=
tscGetTableSchema
(
pTableMetaInfo
->
pTableMeta
);
// functions can not be applied to tags
if
((
index
.
columnIndex
>=
tscGetNumOfColumns
(
pTableMetaInfo
->
pTableMeta
))
||
(
index
.
columnIndex
<
0
))
{
return
invalidSqlErrMsg
(
tscGetErrorMsgPayload
(
pCmd
),
msg6
);
}
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
pSchema
,
cvtFunc
,
pItem
,
colIndex
+
i
,
&
index
,
finalResult
)
!=
0
)
{
char
name
[
TSDB_COL_NAME_LEN
]
=
{
0
};
SSchema
*
pSchema
=
tscGetTableColumnSchema
(
pTableMetaInfo
->
pTableMeta
,
index
.
columnIndex
);
setResultColName
(
name
,
pItem
,
cvtFunc
.
originFuncId
,
&
pParamElem
->
pNode
->
colInfo
);
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
pSchema
,
cvtFunc
,
name
,
colIndex
+
i
,
&
index
,
finalResult
)
!=
0
)
{
return
TSDB_CODE_TSC_INVALID_SQL
;
}
...
...
@@ -2011,7 +2026,12 @@ int32_t addExprAndResultField(SSqlCmd* pCmd, SQueryInfo* pQueryInfo, int32_t col
for
(
int32_t
i
=
0
;
i
<
tscGetNumOfColumns
(
pTableMetaInfo
->
pTableMeta
);
++
i
)
{
SColumnIndex
index
=
{.
tableIndex
=
j
,
.
columnIndex
=
i
};
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
pSchema
,
cvtFunc
,
pItem
,
colIndex
,
&
index
,
finalResult
)
!=
0
)
{
char
name
[
TSDB_COL_NAME_LEN
]
=
{
0
};
SStrToken
t
=
{.
z
=
pSchema
->
name
,
.
n
=
strnlen
(
pSchema
->
name
,
TSDB_COL_NAME_LEN
)};
setResultColName
(
name
,
pItem
,
cvtFunc
.
originFuncId
,
&
t
);
if
(
setExprInfoForFunctions
(
pCmd
,
pQueryInfo
,
&
pSchema
[
index
.
columnIndex
],
cvtFunc
,
name
,
colIndex
,
&
index
,
finalResult
)
!=
0
)
{
return
TSDB_CODE_TSC_INVALID_SQL
;
}
...
...
@@ -3458,7 +3478,7 @@ static int32_t validateArithmeticSQLExpr(SSqlCmd* pCmd, tSQLExpr* pExpr, SQueryI
}
// the expression not from the same table, return error
if
(
uidLeft
!=
uidRight
)
{
if
(
uidLeft
!=
uidRight
&&
uidLeft
!=
0
&&
uidRight
!=
0
)
{
return
TSDB_CODE_TSC_INVALID_SQL
;
}
}
...
...
src/client/src/tscServer.c
浏览文件 @
bc8cffd2
...
...
@@ -880,8 +880,9 @@ int tscBuildQueryMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
pQueryMsg
->
tsOffset
=
htonl
((
int32_t
)(
pMsg
-
pCmd
->
payload
));
if
(
pQueryInfo
->
tsBuf
!=
NULL
)
{
int32_t
vnodeId
=
htonl
(
pQueryMsg
->
head
.
vgId
);
int32_t
code
=
dumpFileBlockByVnodeId
(
pQueryInfo
->
tsBuf
,
vnodeId
,
pMsg
,
&
pQueryMsg
->
tsLen
,
&
pQueryMsg
->
tsNumOfBlocks
);
// note: here used the index instead of actual vnode id.
int32_t
vnodeIndex
=
pTableMetaInfo
->
vgroupIndex
;
int32_t
code
=
dumpFileBlockByVnodeId
(
pQueryInfo
->
tsBuf
,
vnodeIndex
,
pMsg
,
&
pQueryMsg
->
tsLen
,
&
pQueryMsg
->
tsNumOfBlocks
);
if
(
code
!=
TSDB_CODE_SUCCESS
)
{
return
code
;
}
...
...
src/client/src/tscSubquery.c
浏览文件 @
bc8cffd2
...
...
@@ -44,6 +44,17 @@ static int32_t tsCompare(int32_t order, int64_t left, int64_t right) {
}
}
static
void
skipRemainValue
(
STSBuf
*
pTSBuf
,
tVariant
*
tag1
)
{
while
(
tsBufNextPos
(
pTSBuf
))
{
STSElem
el1
=
tsBufGetElem
(
pTSBuf
);
int32_t
res
=
tVariantCompare
(
el1
.
tag
,
tag1
);
if
(
res
!=
0
)
{
// it is a record with new tag
return
;
}
}
}
static
int64_t
doTSBlockIntersect
(
SSqlObj
*
pSql
,
SJoinSupporter
*
pSupporter1
,
SJoinSupporter
*
pSupporter2
,
STimeWindow
*
win
)
{
SQueryInfo
*
pQueryInfo
=
tscGetQueryInfoDetail
(
&
pSql
->
cmd
,
pSql
->
cmd
.
clauseIndex
);
...
...
@@ -92,54 +103,48 @@ static int64_t doTSBlockIntersect(SSqlObj* pSql, SJoinSupporter* pSupporter1, SJ
int64_t
numOfInput1
=
1
;
int64_t
numOfInput2
=
1
;
int32_t
numOfVnodes
=
0
;
int32_t
*
idList
=
NULL
;
tsBufGetVnodeIdList
(
pSupporter2
->
pTSBuf
,
&
numOfVnodes
,
&
idList
);
bool
completed
=
false
;
while
(
1
)
{
STSElem
elem
=
tsBufGetElem
(
pSupporter1
->
pTSBuf
);
// no data in pSupporter1 anymore, jump out of loop
if
(
elem
.
vnode
<
0
||
completed
)
{
if
(
!
tsBufIsValidElem
(
&
elem
)
)
{
break
;
}
bool
f
=
false
;
for
(
int32_t
i
=
0
;
i
<
numOfVnodes
;
++
i
)
{
STSElem
el
=
tsBufGetElemStartPos
(
pSupporter2
->
pTSBuf
,
idList
[
i
],
elem
.
tag
);
if
(
el
.
vnode
==
idList
[
i
])
{
f
=
true
;
break
;
}
}
// find the data in supporter2 with the same tag value
STSElem
e2
=
tsBufFindElemStartPosByTag
(
pSupporter2
->
pTSBuf
,
elem
.
tag
);
/**
* there are elements in pSupporter2 with the same tag, continue
*/
if
(
f
)
{
tVariant
tag1
=
{
0
};
tVariantAssign
(
&
tag1
,
elem
.
tag
);
if
(
tsBufIsValidElem
(
&
e2
))
{
while
(
1
)
{
STSElem
elem1
=
tsBufGetElem
(
pSupporter1
->
pTSBuf
);
STSElem
elem2
=
tsBufGetElem
(
pSupporter2
->
pTSBuf
);
// data with current are exhausted
if
(
!
tsBufIsValidElem
(
&
elem1
)
||
tVariantCompare
(
elem1
.
tag
,
&
tag1
)
!=
0
)
{
break
;
}
if
(
!
tsBufIsValidElem
(
&
elem2
)
||
tVariantCompare
(
elem2
.
tag
,
&
tag1
)
!=
0
)
{
// ignore all records with the same tag
skipRemainValue
(
pSupporter1
->
pTSBuf
,
&
tag1
);
break
;
}
/*
* in case of stable query, limit/offset is not applied here. the limit/offset is applied to the
* final results which is acquired after the secondry merge of in the client.
* final results which is acquired after the second
a
ry merge of in the client.
*/
int32_t
re
=
tsCompare
(
order
,
elem1
.
ts
,
elem2
.
ts
);
if
(
re
<
0
)
{
if
(
!
tsBufNextPos
(
pSupporter1
->
pTSBuf
))
{
completed
=
true
;
break
;
}
tsBufNextPos
(
pSupporter1
->
pTSBuf
);
numOfInput1
++
;
}
else
if
(
re
>
0
)
{
if
(
!
tsBufNextPos
(
pSupporter2
->
pTSBuf
))
{
completed
=
true
;
break
;
}
tsBufNextPos
(
pSupporter2
->
pTSBuf
);
numOfInput2
++
;
}
else
{
if
(
pLimit
->
offset
==
0
||
pQueryInfo
->
interval
.
interval
>
0
||
QUERY_IS_STABLE_QUERY
(
pQueryInfo
->
type
))
{
...
...
@@ -154,43 +159,18 @@ static int64_t doTSBlockIntersect(SSqlObj* pSql, SJoinSupporter* pSupporter1, SJ
tsBufAppend
(
output1
,
elem1
.
vnode
,
elem1
.
tag
,
(
const
char
*
)
&
elem1
.
ts
,
sizeof
(
elem1
.
ts
));
tsBufAppend
(
output2
,
elem2
.
vnode
,
elem2
.
tag
,
(
const
char
*
)
&
elem2
.
ts
,
sizeof
(
elem2
.
ts
));
}
else
{
pLimit
->
offset
-=
1
;
}
if
(
!
tsBufNextPos
(
pSupporter1
->
pTSBuf
))
{
completed
=
true
;
break
;
pLimit
->
offset
-=
1
;
//offset apply to projection?
}
tsBufNextPos
(
pSupporter1
->
pTSBuf
);
numOfInput1
++
;
if
(
!
tsBufNextPos
(
pSupporter2
->
pTSBuf
))
{
completed
=
true
;
break
;
}
tsBufNextPos
(
pSupporter2
->
pTSBuf
);
numOfInput2
++
;
}
}
}
else
{
// no data in pSupporter2, ignore current data in pSupporter2
tVariant
tag
=
{
0
};
tVariantAssign
(
&
tag
,
elem
.
tag
);
// ignore all records with the same tag
while
(
tsBufNextPos
(
pSupporter1
->
pTSBuf
))
{
STSElem
el1
=
tsBufGetElem
(
pSupporter1
->
pTSBuf
);
int32_t
res
=
tVariantCompare
(
el1
.
tag
,
&
tag
);
// it is a record with new tag
if
(
res
!=
0
)
{
break
;
}
}
STSElem
el1
=
tsBufGetElem
(
pSupporter1
->
pTSBuf
);
if
(
el1
.
vnode
<
0
)
{
// no data exists, abort
break
;
}
skipRemainValue
(
pSupporter1
->
pTSBuf
,
&
tag1
);
}
}
...
...
@@ -299,6 +279,68 @@ static UNUSED_FUNC bool needSecondaryQuery(SQueryInfo* pQueryInfo) {
return
false
;
}
static
void
filterVgroupTables
(
SQueryInfo
*
pQueryInfo
,
SArray
*
pVgroupTables
)
{
int32_t
num
=
0
;
int32_t
*
list
=
NULL
;
tsBufGetVnodeIdList
(
pQueryInfo
->
tsBuf
,
&
num
,
&
list
);
// The virtual node, of which all tables are disqualified after the timestamp intersection,
// is removed to avoid next stage query.
// TODO: If tables from some vnodes are not qualified for next stage query, discard them.
for
(
int32_t
k
=
0
;
k
<
taosArrayGetSize
(
pVgroupTables
);)
{
SVgroupTableInfo
*
p
=
taosArrayGet
(
pVgroupTables
,
k
);
bool
found
=
false
;
for
(
int32_t
f
=
0
;
f
<
num
;
++
f
)
{
if
(
p
->
vgInfo
.
vgId
==
list
[
f
])
{
found
=
true
;
break
;
}
}
if
(
!
found
)
{
tscRemoveVgroupTableGroup
(
pVgroupTables
,
k
);
}
else
{
k
++
;
}
}
assert
(
taosArrayGetSize
(
pVgroupTables
)
>
0
);
TSDB_QUERY_SET_TYPE
(
pQueryInfo
->
type
,
TSDB_QUERY_TYPE_MULTITABLE_QUERY
);
taosTFree
(
list
);
}
static
SArray
*
buildVgroupTableByResult
(
SQueryInfo
*
pQueryInfo
,
SArray
*
pVgroupTables
)
{
int32_t
num
=
0
;
int32_t
*
list
=
NULL
;
tsBufGetVnodeIdList
(
pQueryInfo
->
tsBuf
,
&
num
,
&
list
);
int32_t
numOfGroups
=
taosArrayGetSize
(
pVgroupTables
);
SArray
*
pNew
=
taosArrayInit
(
num
,
sizeof
(
SVgroupTableInfo
));
SVgroupTableInfo
info
;
for
(
int32_t
i
=
0
;
i
<
num
;
++
i
)
{
int32_t
vnodeId
=
list
[
i
];
for
(
int32_t
j
=
0
;
j
<
numOfGroups
;
++
j
)
{
SVgroupTableInfo
*
p1
=
taosArrayGet
(
pVgroupTables
,
j
);
if
(
p1
->
vgInfo
.
vgId
==
vnodeId
)
{
tscVgroupTableCopy
(
&
info
,
p1
);
break
;
}
}
taosArrayPush
(
pNew
,
&
info
);
}
taosTFree
(
list
);
TSDB_QUERY_SET_TYPE
(
pQueryInfo
->
type
,
TSDB_QUERY_TYPE_MULTITABLE_QUERY
);
return
pNew
;
}
/*
* launch secondary stage query to fetch the result that contains timestamp in set
*/
...
...
@@ -373,12 +415,11 @@ static int32_t tscLaunchRealSubqueries(SSqlObj* pSql) {
pQueryInfo
->
fieldsInfo
=
pSupporter
->
fieldsInfo
;
pQueryInfo
->
groupbyExpr
=
pSupporter
->
groupInfo
;
SQueryInfo
*
pNewQueryInfo
=
tscGetQueryInfoDetail
(
&
pNew
->
cmd
,
0
);
assert
(
pNew
->
subState
.
numOfSub
==
0
&&
pNew
->
cmd
.
numOfClause
==
1
&&
pNewQueryInfo
->
numOfTables
==
1
);
assert
(
pNew
->
subState
.
numOfSub
==
0
&&
pNew
->
cmd
.
numOfClause
==
1
&&
pQueryInfo
->
numOfTables
==
1
);
tscFieldInfoUpdateOffset
(
p
New
QueryInfo
);
tscFieldInfoUpdateOffset
(
pQueryInfo
);
STableMetaInfo
*
pTableMetaInfo
=
tscGetMetaInfo
(
p
New
QueryInfo
,
0
);
STableMetaInfo
*
pTableMetaInfo
=
tscGetMetaInfo
(
pQueryInfo
,
0
);
pTableMetaInfo
->
pVgroupTables
=
pSupporter
->
pVgroupTables
;
pSupporter
->
exprList
=
NULL
;
...
...
@@ -392,7 +433,7 @@ static int32_t tscLaunchRealSubqueries(SSqlObj* pSql) {
* during the timestamp intersection.
*/
pSupporter
->
limit
=
pQueryInfo
->
limit
;
p
New
QueryInfo
->
limit
=
pSupporter
->
limit
;
pQueryInfo
->
limit
=
pSupporter
->
limit
;
SColumnIndex
index
=
{.
tableIndex
=
0
,
.
columnIndex
=
PRIMARYKEY_TIMESTAMP_COL_INDEX
};
SSchema
*
s
=
tscGetTableColumnSchema
(
pTableMetaInfo
->
pTableMeta
,
0
);
...
...
@@ -407,7 +448,7 @@ static int32_t tscLaunchRealSubqueries(SSqlObj* pSql) {
tscAddSpecialColumnForSelect
(
pQueryInfo
,
0
,
functionId
,
&
index
,
s
,
TSDB_COL_NORMAL
);
tscPrintSelectClause
(
pNew
,
0
);
tscFieldInfoUpdateOffset
(
p
New
QueryInfo
);
tscFieldInfoUpdateOffset
(
pQueryInfo
);
pExpr
=
tscSqlExprGet
(
pQueryInfo
,
0
);
}
...
...
@@ -422,39 +463,21 @@ static int32_t tscLaunchRealSubqueries(SSqlObj* pSql) {
pExpr
->
numOfParams
=
1
;
}
int32_t
num
=
0
;
int32_t
*
list
=
NULL
;
tsBufGetVnodeIdList
(
pNewQueryInfo
->
tsBuf
,
&
num
,
&
list
);
if
(
pTableMetaInfo
->
pVgroupTables
!=
NULL
)
{
for
(
int32_t
k
=
0
;
k
<
taosArrayGetSize
(
pTableMetaInfo
->
pVgroupTables
);)
{
SVgroupTableInfo
*
p
=
taosArrayGet
(
pTableMetaInfo
->
pVgroupTables
,
k
);
bool
found
=
false
;
for
(
int32_t
f
=
0
;
f
<
num
;
++
f
)
{
if
(
p
->
vgInfo
.
vgId
==
list
[
f
])
{
found
=
true
;
break
;
}
}
if
(
!
found
)
{
tscRemoveVgroupTableGroup
(
pTableMetaInfo
->
pVgroupTables
,
k
);
}
else
{
k
++
;
}
if
(
UTIL_TABLE_IS_SUPER_TABLE
(
pTableMetaInfo
))
{
assert
(
pTableMetaInfo
->
pVgroupTables
!=
NULL
);
if
(
tscNonOrderedProjectionQueryOnSTable
(
pQueryInfo
,
0
))
{
SArray
*
p
=
buildVgroupTableByResult
(
pQueryInfo
,
pTableMetaInfo
->
pVgroupTables
);
tscFreeVgroupTableInfo
(
pTableMetaInfo
->
pVgroupTables
);
pTableMetaInfo
->
pVgroupTables
=
p
;
}
else
{
filterVgroupTables
(
pQueryInfo
,
pTableMetaInfo
->
pVgroupTables
);
}
assert
(
taosArrayGetSize
(
pTableMetaInfo
->
pVgroupTables
)
>
0
);
TSDB_QUERY_SET_TYPE
(
pQueryInfo
->
type
,
TSDB_QUERY_TYPE_MULTITABLE_QUERY
);
}
taosTFree
(
list
);
size_t
numOfCols
=
taosArrayGetSize
(
pNewQueryInfo
->
colList
);
size_t
numOfCols
=
taosArrayGetSize
(
pQueryInfo
->
colList
);
tscDebug
(
"%p subquery:%p tableIndex:%d, vgroupIndex:%d, type:%d, exprInfo:%"
PRIzu
", colList:%"
PRIzu
", fieldsInfo:%d, name:%s"
,
pSql
,
pNew
,
0
,
pTableMetaInfo
->
vgroupIndex
,
p
NewQueryInfo
->
type
,
taosArrayGetSize
(
pNew
QueryInfo
->
exprList
),
numOfCols
,
p
New
QueryInfo
->
fieldsInfo
.
numOfOutput
,
pTableMetaInfo
->
name
);
pSql
,
pNew
,
0
,
pTableMetaInfo
->
vgroupIndex
,
p
QueryInfo
->
type
,
taosArrayGetSize
(
p
QueryInfo
->
exprList
),
numOfCols
,
pQueryInfo
->
fieldsInfo
.
numOfOutput
,
pTableMetaInfo
->
name
);
}
//prepare the subqueries object failed, abort
...
...
@@ -1034,6 +1057,7 @@ static void joinRetrieveFinalResCallback(void* param, TAOS_RES* tres, int numOfR
}
}
assert
(
pState
->
numOfRemain
>
0
);
if
(
atomic_sub_fetch_32
(
&
pState
->
numOfRemain
,
1
)
>
0
)
{
tscDebug
(
"%p sub:%p completed, remain:%d, total:%d"
,
pParentSql
,
tres
,
pState
->
numOfRemain
,
pState
->
numOfSub
);
return
;
...
...
@@ -1047,7 +1071,7 @@ static void joinRetrieveFinalResCallback(void* param, TAOS_RES* tres, int numOfR
}
// update the records for each subquery in parent sql object.
bool
isSTableSub
=
tscIsTwoStageSTableQuery
(
pQueryInfo
,
0
);
bool
stableQuery
=
tscIsTwoStageSTableQuery
(
pQueryInfo
,
0
);
for
(
int32_t
i
=
0
;
i
<
pState
->
numOfSub
;
++
i
)
{
if
(
pParentSql
->
pSubs
[
i
]
==
NULL
)
{
tscDebug
(
"%p %p sub:%d not retrieve data"
,
pParentSql
,
NULL
,
i
);
...
...
@@ -1061,7 +1085,7 @@ static void joinRetrieveFinalResCallback(void* param, TAOS_RES* tres, int numOfR
pRes1
->
numOfRows
,
pRes1
->
numOfTotal
);
assert
(
pRes1
->
row
<
pRes1
->
numOfRows
);
}
else
{
if
(
!
isSTableSub
)
{
if
(
!
stableQuery
)
{
pRes1
->
numOfClauseTotal
+=
pRes1
->
numOfRows
;
}
...
...
@@ -1074,7 +1098,7 @@ static void joinRetrieveFinalResCallback(void* param, TAOS_RES* tres, int numOfR
tscBuildResFromSubqueries
(
pParentSql
);
}
void
tscFetchDatablockF
rom
Subquery
(
SSqlObj
*
pSql
)
{
void
tscFetchDatablockF
or
Subquery
(
SSqlObj
*
pSql
)
{
assert
(
pSql
->
subState
.
numOfSub
>=
1
);
int32_t
numOfFetch
=
0
;
...
...
@@ -1136,11 +1160,22 @@ void tscFetchDatablockFromSubquery(SSqlObj* pSql) {
if
(
numOfFetch
<=
0
)
{
bool
tryNextVnode
=
false
;
SSqlObj
*
pp
=
pSql
->
pSubs
[
0
];
SQueryInfo
*
pi
=
tscGetQueryInfoDetail
(
&
pp
->
cmd
,
0
);
bool
orderedPrjQuery
=
false
;
for
(
int32_t
i
=
0
;
i
<
pSql
->
subState
.
numOfSub
;
++
i
)
{
SSqlObj
*
pSub
=
pSql
->
pSubs
[
i
];
if
(
pSub
==
NULL
)
{
continue
;
}
SQueryInfo
*
p
=
tscGetQueryInfoDetail
(
&
pSub
->
cmd
,
0
);
orderedPrjQuery
=
tscNonOrderedProjectionQueryOnSTable
(
p
,
0
);
if
(
orderedPrjQuery
)
{
break
;
}
}
// get the number of subquery that need to retrieve the next vnode.
if
(
tscNonOrderedProjectionQueryOnSTable
(
pi
,
0
)
)
{
if
(
orderedPrjQuery
)
{
for
(
int32_t
i
=
0
;
i
<
pSql
->
subState
.
numOfSub
;
++
i
)
{
SSqlObj
*
pSub
=
pSql
->
pSubs
[
i
];
if
(
pSub
!=
NULL
&&
pSub
->
res
.
row
>=
pSub
->
res
.
numOfRows
&&
pSub
->
res
.
completed
)
{
...
...
@@ -1244,7 +1279,6 @@ void tscSetupOutputColumnIndex(SSqlObj* pSql) {
SSqlCmd
*
pCmd
=
&
pSql
->
cmd
;
SSqlRes
*
pRes
=
&
pSql
->
res
;
tscDebug
(
"%p all subquery response, retrieve data for subclause:%d"
,
pSql
,
pCmd
->
clauseIndex
);
// the column transfer support struct has been built
if
(
pRes
->
pColumnIndex
!=
NULL
)
{
...
...
@@ -1340,21 +1374,23 @@ void tscJoinQueryCallback(void* param, TAOS_RES* tres, int code) {
return
;
}
// wait for the other subqueries response from vnode
if
(
atomic_sub_fetch_32
(
&
pParentSql
->
subState
.
numOfRemain
,
1
)
>
0
)
{
return
;
STableMetaInfo
*
pTableMetaInfo
=
tscGetMetaInfo
(
pQueryInfo
,
0
);
// In case of consequence query from other vnode, do not wait for other query response here.
if
(
!
(
pTableMetaInfo
->
vgroupIndex
>
0
&&
tscNonOrderedProjectionQueryOnSTable
(
pQueryInfo
,
0
)))
{
if
(
atomic_sub_fetch_32
(
&
pParentSql
->
subState
.
numOfRemain
,
1
)
>
0
)
{
return
;
}
}
tscSetupOutputColumnIndex
(
pParentSql
);
STableMetaInfo
*
pTableMetaInfo
=
tscGetMetaInfo
(
pQueryInfo
,
0
);
/**
* if the query is a continue query (vgroupIndex > 0 for projection query) for next vnode, do the retrieval of
* data instead of returning to its invoker
*/
if
(
pTableMetaInfo
->
vgroupIndex
>
0
&&
tscNonOrderedProjectionQueryOnSTable
(
pQueryInfo
,
0
))
{
// pParentSql->subState.numOfRemain = pParentSql->subState.numOfSub; // reset the record value
pSql
->
fp
=
joinRetrieveFinalResCallback
;
// continue retrieve data
pSql
->
cmd
.
command
=
TSDB_SQL_FETCH
;
tscProcessSql
(
pSql
);
...
...
src/client/src/tscUtil.c
浏览文件 @
bc8cffd2
...
...
@@ -1696,6 +1696,17 @@ void tscRemoveVgroupTableGroup(SArray* pVgroupTable, int32_t index) {
taosArrayRemove
(
pVgroupTable
,
index
);
}
void
tscVgroupTableCopy
(
SVgroupTableInfo
*
info
,
SVgroupTableInfo
*
pInfo
)
{
memset
(
info
,
0
,
sizeof
(
SVgroupTableInfo
));
info
->
vgInfo
=
pInfo
->
vgInfo
;
for
(
int32_t
j
=
0
;
j
<
pInfo
->
vgInfo
.
numOfEps
;
++
j
)
{
info
->
vgInfo
.
epAddr
[
j
].
fqdn
=
strdup
(
pInfo
->
vgInfo
.
epAddr
[
j
].
fqdn
);
}
info
->
itemList
=
taosArrayClone
(
pInfo
->
itemList
);
}
SArray
*
tscVgroupTableInfoClone
(
SArray
*
pVgroupTables
)
{
if
(
pVgroupTables
==
NULL
)
{
return
NULL
;
...
...
@@ -1707,14 +1718,8 @@ SArray* tscVgroupTableInfoClone(SArray* pVgroupTables) {
SVgroupTableInfo
info
;
for
(
size_t
i
=
0
;
i
<
num
;
i
++
)
{
SVgroupTableInfo
*
pInfo
=
taosArrayGet
(
pVgroupTables
,
i
);
memset
(
&
info
,
0
,
sizeof
(
SVgroupTableInfo
));
info
.
vgInfo
=
pInfo
->
vgInfo
;
for
(
int32_t
j
=
0
;
j
<
pInfo
->
vgInfo
.
numOfEps
;
++
j
)
{
info
.
vgInfo
.
epAddr
[
j
].
fqdn
=
strdup
(
pInfo
->
vgInfo
.
epAddr
[
j
].
fqdn
);
}
tscVgroupTableCopy
(
&
info
,
pInfo
);
info
.
itemList
=
taosArrayClone
(
pInfo
->
itemList
);
taosArrayPush
(
pa
,
&
info
);
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录