Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
7d50bfcb
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看板
提交
7d50bfcb
编写于
7月 01, 2022
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
more code
上级
36fbf6d7
变更
6
显示空白变更内容
内联
并排
Showing
6 changed file
with
632 addition
and
549 deletion
+632
-549
include/common/tcommon.h
include/common/tcommon.h
+4
-2
source/dnode/vnode/src/inc/tsdb.h
source/dnode/vnode/src/inc/tsdb.h
+20
-24
source/dnode/vnode/src/tsdb/tsdbCommit.c
source/dnode/vnode/src/tsdb/tsdbCommit.c
+30
-35
source/dnode/vnode/src/tsdb/tsdbRead.c
source/dnode/vnode/src/tsdb/tsdbRead.c
+1
-4
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
+483
-228
source/dnode/vnode/src/tsdb/tsdbUtil.c
source/dnode/vnode/src/tsdb/tsdbUtil.c
+94
-256
未找到文件。
include/common/tcommon.h
浏览文件 @
7d50bfcb
...
@@ -61,15 +61,17 @@ typedef struct {
...
@@ -61,15 +61,17 @@ typedef struct {
uint64_t
suid
;
uint64_t
suid
;
}
STableListInfo
;
}
STableListInfo
;
#pragma pack(push, 1)
typedef
struct
SColumnDataAgg
{
typedef
struct
SColumnDataAgg
{
int16_t
colId
;
int16_t
colId
;
int16_t
maxIndex
;
int16_t
minIndex
;
int16_t
minIndex
;
int16_t
maxIndex
;
int16_t
numOfNull
;
int16_t
numOfNull
;
int64_t
sum
;
int64_t
sum
;
int64_t
max
;
int64_t
max
;
int64_t
min
;
int64_t
min
;
}
SColumnDataAgg
;
}
SColumnDataAgg
;
#pragma pack(pop)
typedef
struct
SDataBlockInfo
{
typedef
struct
SDataBlockInfo
{
STimeWindow
window
;
STimeWindow
window
;
...
...
source/dnode/vnode/src/inc/tsdb.h
浏览文件 @
7d50bfcb
...
@@ -119,10 +119,7 @@ int32_t tPutBlockCol(uint8_t *p, void *ph);
...
@@ -119,10 +119,7 @@ int32_t tPutBlockCol(uint8_t *p, void *ph);
int32_t
tGetBlockCol
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tGetBlockCol
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tBlockColCmprFn
(
const
void
*
p1
,
const
void
*
p2
);
int32_t
tBlockColCmprFn
(
const
void
*
p1
,
const
void
*
p2
);
// SBlock
// SBlock
#define tBlockInit() ((SBlock){0})
void
tBlockReset
(
SBlock
*
pBlock
);
void
tBlockReset
(
SBlock
*
pBlock
);
void
tBlockClear
(
SBlock
*
pBlock
);
int32_t
tBlockCopy
(
SBlock
*
pBlockSrc
,
SBlock
*
pBlockDest
);
int32_t
tPutBlock
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tPutBlock
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tGetBlock
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tGetBlock
(
uint8_t
*
p
,
void
*
ph
);
int32_t
tBlockCmprFn
(
const
void
*
p1
,
const
void
*
p2
);
int32_t
tBlockCmprFn
(
const
void
*
p1
,
const
void
*
p2
);
...
@@ -134,11 +131,11 @@ int32_t tGetBlockIdx(uint8_t *p, void *ph);
...
@@ -134,11 +131,11 @@ int32_t tGetBlockIdx(uint8_t *p, void *ph);
int32_t
tCmprBlockIdx
(
void
const
*
lhs
,
void
const
*
rhs
);
int32_t
tCmprBlockIdx
(
void
const
*
lhs
,
void
const
*
rhs
);
// SColdata
// SColdata
#define tColDataInit() ((SColData){0})
#define tColDataInit() ((SColData){0})
void
tColDataReset
(
SColData
*
pColData
,
int16_t
cid
,
int8_t
type
);
void
tColDataReset
(
SColData
*
pColData
,
int16_t
cid
,
int8_t
type
,
int8_t
smaOn
);
void
tColDataClear
(
void
*
ph
);
void
tColDataClear
(
void
*
ph
);
int32_t
tColDataAppendValue
(
SColData
*
pColData
,
SColVal
*
pColVal
);
int32_t
tColDataAppendValue
(
SColData
*
pColData
,
SColVal
*
pColVal
);
int32_t
tColDataCopy
(
SColData
*
pColDataSrc
,
SColData
*
pColDataDest
);
int32_t
tColDataGetValue
(
SColData
*
pColData
,
int32_t
iRow
,
SColVal
*
pColVal
);
int32_t
tColDataGetValue
(
SColData
*
pColData
,
int32_t
iRow
,
SColVal
*
pColVal
);
int32_t
tColDataCopy
(
SColData
*
pColDataSrc
,
SColData
*
pColDataDest
);
// SBlockData
// SBlockData
#define tBlockDataFirstRow(PBLOCKDATA) tsdbRowFromBlockData(PBLOCKDATA, 0)
#define tBlockDataFirstRow(PBLOCKDATA) tsdbRowFromBlockData(PBLOCKDATA, 0)
#define tBlockDataLastRow(PBLOCKDATA) tsdbRowFromBlockData(PBLOCKDATA, (PBLOCKDATA)->nRow - 1)
#define tBlockDataLastRow(PBLOCKDATA) tsdbRowFromBlockData(PBLOCKDATA, (PBLOCKDATA)->nRow - 1)
...
@@ -166,9 +163,8 @@ void tsdbFree(uint8_t *pBuf);
...
@@ -166,9 +163,8 @@ void tsdbFree(uint8_t *pBuf);
#define tMapDataInit() ((SMapData){0})
#define tMapDataInit() ((SMapData){0})
void
tMapDataReset
(
SMapData
*
pMapData
);
void
tMapDataReset
(
SMapData
*
pMapData
);
void
tMapDataClear
(
SMapData
*
pMapData
);
void
tMapDataClear
(
SMapData
*
pMapData
);
int32_t
tMapDataCopy
(
SMapData
*
pMapDataSrc
,
SMapData
*
pMapDataDest
);
int32_t
tMapDataPutItem
(
SMapData
*
pMapData
,
void
*
pItem
,
int32_t
(
*
tPutItemFn
)(
uint8_t
*
,
void
*
));
int32_t
tMapDataPutItem
(
SMapData
*
pMapData
,
void
*
pItem
,
int32_t
(
*
tPutItemFn
)(
uint8_t
*
,
void
*
));
int32_t
tMapDataGetItemByIdx
(
SMapData
*
pMapData
,
int32_t
idx
,
void
*
pItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
));
void
tMapDataGetItemByIdx
(
SMapData
*
pMapData
,
int32_t
idx
,
void
*
pItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
));
int32_t
tMapDataSearch
(
SMapData
*
pMapData
,
void
*
pSearchItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
),
int32_t
tMapDataSearch
(
SMapData
*
pMapData
,
void
*
pSearchItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
),
int32_t
(
*
tItemCmprFn
)(
const
void
*
,
const
void
*
),
void
*
pItem
);
int32_t
(
*
tItemCmprFn
)(
const
void
*
,
const
void
*
),
void
*
pItem
);
int32_t
tPutMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
);
int32_t
tPutMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
);
...
@@ -223,7 +219,6 @@ int32_t tsdbWriteBlockIdx(SDataFWriter *pWriter, SArray *aBlockIdx, uint8_t **pp
...
@@ -223,7 +219,6 @@ int32_t tsdbWriteBlockIdx(SDataFWriter *pWriter, SArray *aBlockIdx, uint8_t **pp
int32_t
tsdbWriteBlock
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
,
SBlockIdx
*
pBlockIdx
);
int32_t
tsdbWriteBlock
(
SDataFWriter
*
pWriter
,
SMapData
*
pMapData
,
uint8_t
**
ppBuf
,
SBlockIdx
*
pBlockIdx
);
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
uint8_t
**
ppBuf2
,
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
uint8_t
**
ppBuf2
,
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int8_t
cmprAlg
);
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int8_t
cmprAlg
);
int32_t
tsdbWriteBlockSMA
(
SDataFWriter
*
pWriter
,
SBlockSMA
*
pBlockSMA
,
int64_t
*
rOffset
,
int64_t
*
rSize
);
SDFileSet
*
tsdbDataFWriterGetWSet
(
SDataFWriter
*
pWriter
);
SDFileSet
*
tsdbDataFWriterGetWSet
(
SDataFWriter
*
pWriter
);
// SDataFReader
// SDataFReader
...
@@ -373,31 +368,32 @@ struct SBlockIdx {
...
@@ -373,31 +368,32 @@ struct SBlockIdx {
struct
SMapData
{
struct
SMapData
{
int32_t
nItem
;
int32_t
nItem
;
uint8_t
flag
;
int32_t
*
aOffset
;
uint8_t
*
pOfst
;
int32_t
nData
;
uint32_t
nData
;
uint8_t
*
pData
;
uint8_t
*
pData
;
uint8_t
*
pBuf
;
};
};
typedef
struct
{
typedef
struct
{
int16_t
cid
;
int16_t
cid
;
int8_t
type
;
int8_t
type
;
int8_t
flag
;
// HAS_NONE|HAS_NULL|HAS_VALUE
int8_t
flag
;
// HAS_NONE|HAS_NULL|HAS_VALUE
int64_t
offset
;
int32_t
offset
;
int64_t
bsize
;
// bitmap size
int32_t
szBitmap
;
// bitmap size
int64_t
csize
;
// compressed column value size
int32_t
szOffset
;
// size of offset, only for variant-length data type
int64_t
osize
;
// original column value size (only save for variant data type)
int32_t
szValue
;
// compressed column value size
int32_t
szOrigin
;
// original column value size (only save for variant data type)
}
SBlockCol
;
}
SBlockCol
;
typedef
struct
{
typedef
struct
{
int
64_t
nRow
;
int
32_t
nRow
;
int8_t
cmprAlg
;
int8_t
cmprAlg
;
int64_t
offset
;
int64_t
offset
;
// block data offset
int64_t
szVersion
;
// VERSION size
int32_t
szBlockCol
;
// SBlockCol size
int64_t
szTSKEY
;
// TSKEY size
int32_t
szVersion
;
// VERSION size
int64_t
szBlock
;
// total block size
int32_t
szTSKEY
;
// TSKEY size
SMapData
mBlockCol
;
// SMapData<SBlockCol>
int32_t
szBlock
;
// total block size
int64_t
sOffset
;
// sma offset
int32_t
nSma
;
// sma size
}
SSubBlock
;
}
SSubBlock
;
struct
SBlock
{
struct
SBlock
{
...
@@ -425,7 +421,7 @@ struct SAggrBlkCol {
...
@@ -425,7 +421,7 @@ struct SAggrBlkCol {
struct
SColData
{
struct
SColData
{
int16_t
cid
;
int16_t
cid
;
int8_t
type
;
int8_t
type
;
int8_t
offsetValid
;
int8_t
smaOn
;
int32_t
nVal
;
int32_t
nVal
;
uint8_t
flag
;
uint8_t
flag
;
uint8_t
*
pBitMap
;
uint8_t
*
pBitMap
;
...
...
source/dnode/vnode/src/tsdb/tsdbCommit.c
浏览文件 @
7d50bfcb
...
@@ -33,12 +33,10 @@ typedef struct {
...
@@ -33,12 +33,10 @@ typedef struct {
SDataFReader
*
pReader
;
SDataFReader
*
pReader
;
SArray
*
aBlockIdx
;
// SArray<SBlockIdx>
SArray
*
aBlockIdx
;
// SArray<SBlockIdx>
SMapData
oBlockMap
;
// SMapData<SBlock>, read from reader
SMapData
oBlockMap
;
// SMapData<SBlock>, read from reader
SBlock
oBlock
;
SBlockData
oBlockData
;
SBlockData
oBlockData
;
SDataFWriter
*
pWriter
;
SDataFWriter
*
pWriter
;
SArray
*
aBlockIdxN
;
// SArray<SBlockIdx>
SArray
*
aBlockIdxN
;
// SArray<SBlockIdx>
SMapData
nBlockMap
;
// SMapData<SBlock>
SMapData
nBlockMap
;
// SMapData<SBlock>
SBlock
nBlock
;
SBlockData
nBlockData
;
SBlockData
nBlockData
;
int64_t
suid
;
int64_t
suid
;
int64_t
uid
;
int64_t
uid
;
...
@@ -260,7 +258,6 @@ static int32_t tsdbCommitFileDataStart(SCommitter *pCommitter) {
...
@@ -260,7 +258,6 @@ static int32_t tsdbCommitFileDataStart(SCommitter *pCommitter) {
// old
// old
taosArrayClear
(
pCommitter
->
aBlockIdx
);
taosArrayClear
(
pCommitter
->
aBlockIdx
);
tMapDataReset
(
&
pCommitter
->
oBlockMap
);
tMapDataReset
(
&
pCommitter
->
oBlockMap
);
tBlockReset
(
&
pCommitter
->
oBlock
);
tBlockDataReset
(
&
pCommitter
->
oBlockData
);
tBlockDataReset
(
&
pCommitter
->
oBlockData
);
pRSet
=
tsdbFSStateGetDFileSet
(
pTsdb
->
fs
->
nState
,
pCommitter
->
commitFid
);
pRSet
=
tsdbFSStateGetDFileSet
(
pTsdb
->
fs
->
nState
,
pCommitter
->
commitFid
);
if
(
pRSet
)
{
if
(
pRSet
)
{
...
@@ -274,7 +271,6 @@ static int32_t tsdbCommitFileDataStart(SCommitter *pCommitter) {
...
@@ -274,7 +271,6 @@ static int32_t tsdbCommitFileDataStart(SCommitter *pCommitter) {
// new
// new
taosArrayClear
(
pCommitter
->
aBlockIdxN
);
taosArrayClear
(
pCommitter
->
aBlockIdxN
);
tMapDataReset
(
&
pCommitter
->
nBlockMap
);
tMapDataReset
(
&
pCommitter
->
nBlockMap
);
tBlockReset
(
&
pCommitter
->
nBlock
);
tBlockDataReset
(
&
pCommitter
->
nBlockData
);
tBlockDataReset
(
&
pCommitter
->
nBlockData
);
if
(
pRSet
)
{
if
(
pRSet
)
{
wSet
=
(
SDFileSet
){.
diskId
=
pRSet
->
diskId
,
wSet
=
(
SDFileSet
){.
diskId
=
pRSet
->
diskId
,
...
@@ -351,11 +347,6 @@ static int32_t tsdbCommitBlockData(SCommitter *pCommitter, SBlockData *pBlockDat
...
@@ -351,11 +347,6 @@ static int32_t tsdbCommitBlockData(SCommitter *pCommitter, SBlockData *pBlockDat
code
=
tsdbWriteBlockData
(
pCommitter
->
pWriter
,
pBlockData
,
NULL
,
NULL
,
pBlockIdx
,
pBlock
,
pCommitter
->
cmprAlg
);
code
=
tsdbWriteBlockData
(
pCommitter
->
pWriter
,
pBlockData
,
NULL
,
NULL
,
pBlockIdx
,
pBlock
,
pCommitter
->
cmprAlg
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
#if 0
code = tsdbWriteBlockSMA(pCommitter, pBlockData, pBlock);
if (code) goto _err;
#endif
code
=
tMapDataPutItem
(
&
pCommitter
->
nBlockMap
,
pBlock
,
tPutBlock
);
code
=
tMapDataPutItem
(
&
pCommitter
->
nBlockMap
,
pBlock
,
tPutBlock
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
@@ -371,7 +362,8 @@ static int32_t tsdbMergeTableData(SCommitter *pCommitter, STbDataIter *pIter, SB
...
@@ -371,7 +362,8 @@ static int32_t tsdbMergeTableData(SCommitter *pCommitter, STbDataIter *pIter, SB
SBlockIdx
*
pBlockIdx
=
&
(
SBlockIdx
){.
suid
=
pIter
->
pTbData
->
suid
,
.
uid
=
pIter
->
pTbData
->
uid
};
SBlockIdx
*
pBlockIdx
=
&
(
SBlockIdx
){.
suid
=
pIter
->
pTbData
->
suid
,
.
uid
=
pIter
->
pTbData
->
uid
};
SBlockData
*
pBlockDataMerge
=
&
pCommitter
->
oBlockData
;
SBlockData
*
pBlockDataMerge
=
&
pCommitter
->
oBlockData
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
SBlock
*
pBlock
=
&
pCommitter
->
nBlock
;
SBlock
block
;
SBlock
*
pBlock
=
&
block
;
TSDBROW
*
pRow1
;
TSDBROW
*
pRow1
;
TSDBROW
row2
;
TSDBROW
row2
;
TSDBROW
*
pRow2
=
&
row2
;
TSDBROW
*
pRow2
=
&
row2
;
...
@@ -469,7 +461,8 @@ _err:
...
@@ -469,7 +461,8 @@ _err:
static
int32_t
tsdbCommitTableMemData
(
SCommitter
*
pCommitter
,
STbDataIter
*
pIter
,
TSDBKEY
toKey
,
int8_t
toDataOnly
)
{
static
int32_t
tsdbCommitTableMemData
(
SCommitter
*
pCommitter
,
STbDataIter
*
pIter
,
TSDBKEY
toKey
,
int8_t
toDataOnly
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
TSDBROW
*
pRow
;
TSDBROW
*
pRow
;
SBlock
*
pBlock
=
&
pCommitter
->
nBlock
;
SBlock
block
;
SBlock
*
pBlock
=
&
block
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
int64_t
suid
=
pIter
->
pTbData
->
suid
;
int64_t
suid
=
pIter
->
pTbData
->
suid
;
int64_t
uid
=
pIter
->
pTbData
->
uid
;
int64_t
uid
=
pIter
->
pTbData
->
uid
;
...
@@ -519,13 +512,14 @@ _err:
...
@@ -519,13 +512,14 @@ _err:
static
int32_t
tsdbCommitTableDiskData
(
SCommitter
*
pCommitter
,
SBlock
*
pBlock
,
SBlockIdx
*
pBlockIdx
)
{
static
int32_t
tsdbCommitTableDiskData
(
SCommitter
*
pCommitter
,
SBlock
*
pBlock
,
SBlockIdx
*
pBlockIdx
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
SBlock
block
;
if
(
pBlock
->
last
)
{
if
(
pBlock
->
last
)
{
code
=
tsdbReadBlockData
(
pCommitter
->
pReader
,
pBlockIdx
,
pBlock
,
&
pCommitter
->
oBlockData
,
NULL
,
NULL
);
code
=
tsdbReadBlockData
(
pCommitter
->
pReader
,
pBlockIdx
,
pBlock
,
&
pCommitter
->
oBlockData
,
NULL
,
NULL
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
tBlockReset
(
&
pCommitter
->
nB
lock
);
tBlockReset
(
&
b
lock
);
code
=
tsdbCommitBlockData
(
pCommitter
,
&
pCommitter
->
oBlockData
,
&
pCommitter
->
nB
lock
,
pBlockIdx
,
0
);
code
=
tsdbCommitBlockData
(
pCommitter
,
&
pCommitter
->
oBlockData
,
&
b
lock
,
pBlockIdx
,
0
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
else
{
}
else
{
code
=
tMapDataPutItem
(
&
pCommitter
->
nBlockMap
,
pBlock
,
tPutBlock
);
code
=
tMapDataPutItem
(
&
pCommitter
->
nBlockMap
,
pBlock
,
tPutBlock
);
...
@@ -590,6 +584,7 @@ static int32_t tsdbMergeAsSubBlock(SCommitter *pCommitter, STbDataIter *pIter, S
...
@@ -590,6 +584,7 @@ static int32_t tsdbMergeAsSubBlock(SCommitter *pCommitter, STbDataIter *pIter, S
int32_t
code
=
0
;
int32_t
code
=
0
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
SBlockData
*
pBlockData
=
&
pCommitter
->
nBlockData
;
SBlockIdx
*
pBlockIdx
=
&
(
SBlockIdx
){.
suid
=
pIter
->
pTbData
->
suid
,
.
uid
=
pIter
->
pTbData
->
uid
};
SBlockIdx
*
pBlockIdx
=
&
(
SBlockIdx
){.
suid
=
pIter
->
pTbData
->
suid
,
.
uid
=
pIter
->
pTbData
->
uid
};
SBlock
block
;
TSDBROW
*
pRow
;
TSDBROW
*
pRow
;
tBlockDataReset
(
pBlockData
);
tBlockDataReset
(
pBlockData
);
...
@@ -617,11 +612,8 @@ static int32_t tsdbMergeAsSubBlock(SCommitter *pCommitter, STbDataIter *pIter, S
...
@@ -617,11 +612,8 @@ static int32_t tsdbMergeAsSubBlock(SCommitter *pCommitter, STbDataIter *pIter, S
}
}
}
}
// write as a subblock
block
=
*
pBlock
;
code
=
tBlockCopy
(
pBlock
,
&
pCommitter
->
nBlock
);
code
=
tsdbCommitBlockData
(
pCommitter
,
pBlockData
,
&
block
,
pBlockIdx
,
0
);
if
(
code
)
goto
_err
;
code
=
tsdbCommitBlockData
(
pCommitter
,
pBlockData
,
&
pCommitter
->
nBlock
,
pBlockIdx
,
0
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
return
code
;
return
code
;
...
@@ -670,7 +662,8 @@ static int32_t tsdbCommitTableData(SCommitter *pCommitter, STbData *pTbData, SBl
...
@@ -670,7 +662,8 @@ static int32_t tsdbCommitTableData(SCommitter *pCommitter, STbData *pTbData, SBl
// start ===========
// start ===========
tMapDataReset
(
&
pCommitter
->
nBlockMap
);
tMapDataReset
(
&
pCommitter
->
nBlockMap
);
SBlock
*
pBlock
=
&
pCommitter
->
oBlock
;
SBlock
block
;
SBlock
*
pBlock
=
&
block
;
iBlock
=
0
;
iBlock
=
0
;
if
(
iBlock
<
nBlock
)
{
if
(
iBlock
<
nBlock
)
{
...
@@ -895,6 +888,8 @@ static int32_t tsdbCommitFileData(SCommitter *pCommitter) {
...
@@ -895,6 +888,8 @@ static int32_t tsdbCommitFileData(SCommitter *pCommitter) {
_err:
_err:
tsdbError
(
"vgId:%d commit file data failed since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbError
(
"vgId:%d commit file data failed since %s"
,
TD_VID
(
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbDataFReaderClose
(
&
pCommitter
->
pReader
);
tsdbDataFWriterClose
(
&
pCommitter
->
pWriter
,
0
);
return
code
;
return
code
;
}
}
...
@@ -931,21 +926,23 @@ _err:
...
@@ -931,21 +926,23 @@ _err:
static
int32_t
tsdbCommitDataStart
(
SCommitter
*
pCommitter
)
{
static
int32_t
tsdbCommitDataStart
(
SCommitter
*
pCommitter
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
pCommitter
->
pReader
=
NULL
;
pCommitter
->
aBlockIdx
=
taosArrayInit
(
0
,
sizeof
(
SBlockIdx
));
pCommitter
->
aBlockIdx
=
taosArrayInit
(
0
,
sizeof
(
SBlockIdx
));
pCommitter
->
oBlockMap
=
tMapDataInit
();
if
(
pCommitter
->
aBlockIdx
==
NULL
)
{
pCommitter
->
oBlock
=
tBlockInit
();
code
=
TSDB_CODE_OUT_OF_MEMORY
;
pCommitter
->
pWriter
=
NULL
;
goto
_exit
;
}
pCommitter
->
aBlockIdxN
=
taosArrayInit
(
0
,
sizeof
(
SBlockIdx
));
pCommitter
->
aBlockIdxN
=
taosArrayInit
(
0
,
sizeof
(
SBlockIdx
));
pCommitter
->
nBlockMap
=
tMapDataInit
();
if
(
pCommitter
->
aBlockIdxN
==
NULL
)
{
pCommitter
->
nBlock
=
tBlockInit
();
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_exit
;
}
code
=
tBlockDataInit
(
&
pCommitter
->
oBlockData
);
code
=
tBlockDataInit
(
&
pCommitter
->
oBlockData
);
if
(
code
)
goto
_exit
;
if
(
code
)
goto
_exit
;
code
=
tBlockDataInit
(
&
pCommitter
->
nBlockData
);
code
=
tBlockDataInit
(
&
pCommitter
->
nBlockData
);
if
(
code
)
{
if
(
code
)
goto
_exit
;
tBlockDataClear
(
&
pCommitter
->
oBlockData
);
goto
_exit
;
}
_exit:
_exit:
return
code
;
return
code
;
...
@@ -953,13 +950,11 @@ _exit:
...
@@ -953,13 +950,11 @@ _exit:
static
void
tsdbCommitDataEnd
(
SCommitter
*
pCommitter
)
{
static
void
tsdbCommitDataEnd
(
SCommitter
*
pCommitter
)
{
taosArrayDestroy
(
pCommitter
->
aBlockIdx
);
taosArrayDestroy
(
pCommitter
->
aBlockIdx
);
// tMapDataClear(&pCommitter->oBlockMap);
tMapDataClear
(
&
pCommitter
->
oBlockMap
);
// tBlockClear(&pCommitter->oBlock);
tBlockDataClear
(
&
pCommitter
->
oBlockData
);
// tBlockDataClear(&pCommitter->oBlockData);
taosArrayDestroy
(
pCommitter
->
aBlockIdxN
);
taosArrayDestroy
(
pCommitter
->
aBlockIdxN
);
// tMapDataClear(&pCommitter->nBlockMap);
tMapDataClear
(
&
pCommitter
->
nBlockMap
);
// tBlockClear(&pCommitter->nBlock);
tBlockDataClear
(
&
pCommitter
->
nBlockData
);
// tBlockDataClear(&pCommitter->nBlockData);
}
}
static
int32_t
tsdbCommitData
(
SCommitter
*
pCommitter
)
{
static
int32_t
tsdbCommitData
(
SCommitter
*
pCommitter
)
{
...
...
source/dnode/vnode/src/tsdb/tsdbRead.c
浏览文件 @
7d50bfcb
...
@@ -807,10 +807,7 @@ static int32_t doLoadFileBlock(STsdbReader* pReader, SArray* pIndexList, uint32_
...
@@ -807,10 +807,7 @@ static int32_t doLoadFileBlock(STsdbReader* pReader, SArray* pIndexList, uint32_
for
(
int32_t
j
=
0
;
j
<
mapData
.
nItem
;
++
j
)
{
for
(
int32_t
j
=
0
;
j
<
mapData
.
nItem
;
++
j
)
{
SBlock
block
=
{
0
};
SBlock
block
=
{
0
};
int32_t
code
=
tMapDataGetItemByIdx
(
&
mapData
,
j
,
&
block
,
tGetBlock
);
tMapDataGetItemByIdx
(
&
mapData
,
j
,
&
block
,
tGetBlock
);
if
(
code
!=
TSDB_CODE_SUCCESS
)
{
return
code
;
}
// 1. time range check
// 1. time range check
if
(
block
.
minKey
.
ts
>
pReader
->
window
.
ekey
||
block
.
maxKey
.
ts
<
pReader
->
window
.
skey
)
{
if
(
block
.
minKey
.
ts
>
pReader
->
window
.
ekey
||
block
.
maxKey
.
ts
<
pReader
->
window
.
skey
)
{
...
...
source/dnode/vnode/src/tsdb/tsdbReaderWriter.c
浏览文件 @
7d50bfcb
...
@@ -489,6 +489,7 @@ _err:
...
@@ -489,6 +489,7 @@ _err:
int32_t
tsdbDataFReaderClose
(
SDataFReader
**
ppReader
)
{
int32_t
tsdbDataFReaderClose
(
SDataFReader
**
ppReader
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
if
(
*
ppReader
==
NULL
)
goto
_exit
;
if
(
taosCloseFile
(
&
(
*
ppReader
)
->
pHeadFD
)
<
0
)
{
if
(
taosCloseFile
(
&
(
*
ppReader
)
->
pHeadFD
)
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
...
@@ -511,6 +512,8 @@ int32_t tsdbDataFReaderClose(SDataFReader **ppReader) {
...
@@ -511,6 +512,8 @@ int32_t tsdbDataFReaderClose(SDataFReader **ppReader) {
}
}
taosMemoryFree
(
*
ppReader
);
taosMemoryFree
(
*
ppReader
);
_exit:
*
ppReader
=
NULL
;
*
ppReader
=
NULL
;
return
code
;
return
code
;
...
@@ -586,11 +589,14 @@ int32_t tsdbReadBlock(SDataFReader *pReader, SBlockIdx *pBlockIdx, SMapData *mBl
...
@@ -586,11 +589,14 @@ int32_t tsdbReadBlock(SDataFReader *pReader, SBlockIdx *pBlockIdx, SMapData *mBl
int32_t
code
=
0
;
int32_t
code
=
0
;
int64_t
offset
=
pBlockIdx
->
offset
;
int64_t
offset
=
pBlockIdx
->
offset
;
int64_t
size
=
pBlockIdx
->
size
;
int64_t
size
=
pBlockIdx
->
size
;
uint8_t
*
pBuf
=
NULL
;
int64_t
n
;
int64_t
n
;
int64_t
tn
;
SBlockDataHdr
hdr
;
SBlockDataHdr
hdr
;
if
(
!
ppBuf
)
ppBuf
=
&
pBuf
;
// alloc
// alloc
if
(
!
ppBuf
)
ppBuf
=
&
mBlock
->
pBuf
;
code
=
tsdbRealloc
(
ppBuf
,
size
);
code
=
tsdbRealloc
(
ppBuf
,
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
@@ -623,17 +629,24 @@ int32_t tsdbReadBlock(SDataFReader *pReader, SBlockIdx *pBlockIdx, SMapData *mBl
...
@@ -623,17 +629,24 @@ int32_t tsdbReadBlock(SDataFReader *pReader, SBlockIdx *pBlockIdx, SMapData *mBl
ASSERT
(
hdr
.
uid
==
pBlockIdx
->
uid
);
ASSERT
(
hdr
.
uid
==
pBlockIdx
->
uid
);
n
=
sizeof
(
hdr
);
n
=
sizeof
(
hdr
);
n
+=
tGetMapData
(
*
ppBuf
+
n
,
mBlock
);
tn
=
tGetMapData
(
*
ppBuf
+
n
,
mBlock
);
if
(
tn
<
0
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
n
+=
tn
;
ASSERT
(
n
+
sizeof
(
TSCKSUM
)
==
size
);
ASSERT
(
n
+
sizeof
(
TSCKSUM
)
==
size
);
tsdbFree
(
pBuf
);
return
code
;
return
code
;
_err:
_err:
tsdbError
(
"vgId:%d read block failed since %s"
,
TD_VID
(
pReader
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbError
(
"vgId:%d read block failed since %s"
,
TD_VID
(
pReader
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbFree
(
pBuf
);
return
code
;
return
code
;
}
}
static
int32_t
tsdbRe
cover
BlockDataKey
(
SBlockData
*
pBlockData
,
SSubBlock
*
pSubBlock
,
uint8_t
*
pBuf
,
uint8_t
**
ppBuf
)
{
static
int32_t
tsdbRe
ad
BlockDataKey
(
SBlockData
*
pBlockData
,
SSubBlock
*
pSubBlock
,
uint8_t
*
pBuf
,
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
int64_t
size
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
int64_t
size
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
int64_t
n
;
int64_t
n
;
...
@@ -656,7 +669,6 @@ static int32_t tsdbRecoverBlockDataKey(SBlockData *pBlockData, SSubBlock *pSubBl
...
@@ -656,7 +669,6 @@ static int32_t tsdbRecoverBlockDataKey(SBlockData *pBlockData, SSubBlock *pSubBl
memcpy
(
pBlockData
->
aVersion
,
pBuf
,
pSubBlock
->
szVersion
);
memcpy
(
pBlockData
->
aVersion
,
pBuf
,
pSubBlock
->
szVersion
);
// TSKEY
// TSKEY
pBuf
=
pBuf
+
pSubBlock
->
szVersion
;
memcpy
(
pBlockData
->
aTSKEY
,
pBuf
+
pSubBlock
->
szVersion
,
pSubBlock
->
szTSKEY
);
memcpy
(
pBlockData
->
aTSKEY
,
pBuf
+
pSubBlock
->
szVersion
,
pSubBlock
->
szTSKEY
);
}
else
{
}
else
{
size
=
sizeof
(
int64_t
)
*
pSubBlock
->
nRow
+
COMP_OVERFLOW_BYTES
;
size
=
sizeof
(
int64_t
)
*
pSubBlock
->
nRow
+
COMP_OVERFLOW_BYTES
;
...
@@ -674,9 +686,9 @@ static int32_t tsdbRecoverBlockDataKey(SBlockData *pBlockData, SSubBlock *pSubBl
...
@@ -674,9 +686,9 @@ static int32_t tsdbRecoverBlockDataKey(SBlockData *pBlockData, SSubBlock *pSubBl
}
}
// TSKEY
// TSKEY
pBuf
=
pBuf
+
pSubBlock
->
szVersion
;
n
=
tsDecompressTimestamp
(
pBuf
+
pSubBlock
->
szVersion
,
pSubBlock
->
szTSKEY
,
pSubBlock
->
nRow
,
n
=
tsDecompressTimestamp
(
pBuf
,
pSubBlock
->
szTSKEY
,
pSubBlock
->
nRow
,
(
char
*
)
pBlockData
->
aTSKEY
,
(
char
*
)
pBlockData
->
aTSKEY
,
sizeof
(
TSKEY
)
*
pSubBlock
->
nRow
,
pSubBlock
->
cmprAlg
,
*
ppBuf
,
size
of
(
TSKEY
)
*
pSubBlock
->
nRow
,
pSubBlock
->
cmprAlg
,
*
ppBuf
,
size
);
size
);
if
(
n
<
0
)
{
if
(
n
<
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
...
@@ -689,15 +701,13 @@ _err:
...
@@ -689,15 +701,13 @@ _err:
return
code
;
return
code
;
}
}
static
int32_t
tsdbRe
coverColData
(
SBlockData
*
pBlockData
,
SSubBlock
*
pSubBlock
,
SBlockCol
*
pBlockCol
,
static
int32_t
tsdbRe
adColDataImpl
(
SSubBlock
*
pSubBlock
,
SBlockCol
*
pBlockCol
,
SColData
*
pColData
,
uint8_t
*
pBuf
,
SColData
*
pColData
,
uint8_t
*
pBuf
,
uint8_t
**
ppBuf
)
{
uint8_t
**
ppBuf
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
int64_t
size
;
int64_t
size
;
int64_t
n
;
int64_t
n
;
ASSERT
(
pBlockCol
->
flag
!=
HAS_NULL
);
if
(
!
taosCheckChecksumWhole
(
pBuf
,
pBlockCol
->
szBitmap
+
pBlockCol
->
szOffset
+
pBlockCol
->
szValue
+
sizeof
(
TSCKSUM
)))
{
if
(
!
taosCheckChecksumWhole
(
pBuf
,
pBlockCol
->
bsize
+
pBlockCol
->
csize
+
sizeof
(
TSCKSUM
)))
{
code
=
TSDB_CODE_FILE_CORRUPTED
;
code
=
TSDB_CODE_FILE_CORRUPTED
;
goto
_err
;
goto
_err
;
}
}
...
@@ -705,40 +715,71 @@ static int32_t tsdbRecoverColData(SBlockData *pBlockData, SSubBlock *pSubBlock,
...
@@ -705,40 +715,71 @@ static int32_t tsdbRecoverColData(SBlockData *pBlockData, SSubBlock *pSubBlock,
pColData
->
nVal
=
pSubBlock
->
nRow
;
pColData
->
nVal
=
pSubBlock
->
nRow
;
pColData
->
flag
=
pBlockCol
->
flag
;
pColData
->
flag
=
pBlockCol
->
flag
;
//
bitmap
//
BITMAP
if
(
pBlockCol
->
flag
!=
HAS_VALUE
)
{
if
(
pBlockCol
->
flag
!=
HAS_VALUE
)
{
size
=
BIT2_SIZE
(
pSubBlock
->
nRow
);
ASSERT
(
pBlockCol
->
szBitmap
);
size
=
BIT2_SIZE
(
pColData
->
nVal
);
code
=
tsdbRealloc
(
&
pColData
->
pBitMap
,
size
);
code
=
tsdbRealloc
(
&
pColData
->
pBitMap
,
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
ASSERT
(
pBlockCol
->
bsize
==
size
);
code
=
tsdbRealloc
(
ppBuf
,
size
+
COMP_OVERFLOW_BYTES
);
if
(
code
)
goto
_err
;
n
=
tsDecompressTinyint
(
pBuf
,
pBlockCol
->
szBitmap
,
size
,
pColData
->
pBitMap
,
size
,
TWO_STAGE_COMP
,
*
ppBuf
,
size
+
COMP_OVERFLOW_BYTES
);
if
(
n
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
}
memcpy
(
pColData
->
pBitMap
,
pBuf
,
size
);
ASSERT
(
n
==
size
);
}
else
{
}
else
{
ASSERT
(
pBlockCol
->
bsize
==
0
);
ASSERT
(
pBlockCol
->
szBitmap
==
0
);
}
}
pBuf
=
pBuf
+
pBlockCol
->
bsize
;
pBuf
=
pBuf
+
pBlockCol
->
szBitmap
;
// value
// OFFSET
if
(
IS_VAR_DATA_TYPE
(
pBlockCol
->
type
))
{
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
{
pColData
->
nData
=
pBlockCol
->
osize
;
ASSERT
(
pBlockCol
->
szOffset
);
size
=
sizeof
(
int32_t
)
*
pColData
->
nVal
;
code
=
tsdbRealloc
((
uint8_t
**
)
&
pColData
->
aOffset
,
size
);
if
(
code
)
goto
_err
;
code
=
tsdbRealloc
(
ppBuf
,
size
+
COMP_OVERFLOW_BYTES
);
if
(
code
)
goto
_err
;
n
=
tsDecompressInt
(
pBuf
,
pBlockCol
->
szOffset
,
pColData
->
nVal
,
(
char
*
)
pColData
->
aOffset
,
size
,
TWO_STAGE_COMP
,
*
ppBuf
,
size
+
COMP_OVERFLOW_BYTES
);
if
(
n
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
}
ASSERT
(
n
==
size
);
}
else
{
}
else
{
pColData
->
nData
=
tDataTypes
[
pBlockCol
->
type
].
bytes
*
pSubBlock
->
nRow
;
ASSERT
(
pBlockCol
->
szOffset
==
0
)
;
}
}
pBuf
=
pBuf
+
pBlockCol
->
szOffset
;
// VALUE
pColData
->
nData
=
pBlockCol
->
szOrigin
;
code
=
tsdbRealloc
(
&
pColData
->
pData
,
pColData
->
nData
);
code
=
tsdbRealloc
(
&
pColData
->
pData
,
pColData
->
nData
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
if
(
pSubBlock
->
cmprAlg
==
NO_COMPRESSION
)
{
if
(
pSubBlock
->
cmprAlg
==
NO_COMPRESSION
)
{
memcpy
(
pColData
->
pData
,
pBuf
,
pColData
->
nData
);
memcpy
(
pColData
->
pData
,
pBuf
,
pColData
->
nData
);
}
else
{
}
else
{
size
=
pColData
->
nData
+
COMP_OVERFLOW_BYTES
;
if
(
pSubBlock
->
cmprAlg
==
TWO_STAGE_COMP
)
{
if
(
pSubBlock
->
cmprAlg
==
TWO_STAGE_COMP
)
{
code
=
tsdbRealloc
(
ppBuf
,
size
);
code
=
tsdbRealloc
(
ppBuf
,
pColData
->
nData
+
COMP_OVERFLOW_BYTES
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
}
n
=
tDataTypes
[
pBlockCol
->
type
].
decompFunc
(
pBuf
,
pBlockCol
->
csize
,
pSubBlock
->
nRow
,
pColData
->
pData
,
n
=
tDataTypes
[
pBlockCol
->
type
].
decompFunc
(
pBuf
,
pBlockCol
->
szValue
,
pSubBlock
->
nRow
,
pColData
->
pData
,
pColData
->
nData
,
pSubBlock
->
cmprAlg
,
*
ppBuf
,
size
);
pColData
->
nData
,
pSubBlock
->
cmprAlg
,
*
ppBuf
,
pColData
->
nData
+
COMP_OVERFLOW_BYTES
);
if
(
n
<
0
)
{
if
(
n
<
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
...
@@ -753,11 +794,41 @@ _err:
...
@@ -753,11 +794,41 @@ _err:
return
code
;
return
code
;
}
}
static
int32_t
tsdbReadColDataImpl
(
SDataFReader
*
pReader
,
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int32_t
iSubBlock
,
static
int32_t
tsdbReadBlockCol
(
SSubBlock
*
pSubBlock
,
uint8_t
*
p
,
SArray
*
aBlockCol
)
{
int32_t
code
=
0
;
int32_t
n
=
0
;
SBlockCol
blockCol
;
SBlockCol
*
pBlockCol
=
&
blockCol
;
if
(
!
taosCheckChecksumWhole
(
p
,
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
)))
{
code
=
TSDB_CODE_FILE_CORRUPTED
;
goto
_err
;
}
n
+=
sizeof
(
SBlockDataHdr
);
while
(
n
<
pSubBlock
->
szBlockCol
)
{
n
+=
tGetBlockCol
(
p
+
n
,
pBlockCol
);
if
(
taosArrayPush
(
aBlockCol
,
pBlockCol
)
==
NULL
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
}
ASSERT
(
n
==
pSubBlock
->
szBlockCol
);
return
code
;
_err:
return
code
;
}
static
int32_t
tsdbReadSubColData
(
SDataFReader
*
pReader
,
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int32_t
iSubBlock
,
int16_t
*
aColId
,
int32_t
nCol
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
int16_t
*
aColId
,
int32_t
nCol
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
uint8_t
**
ppBuf2
)
{
uint8_t
**
ppBuf2
)
{
TdFilePtr
pFD
=
pBlock
->
last
?
pReader
->
pLastFD
:
pReader
->
pDataFD
;
TdFilePtr
pFD
=
pBlock
->
last
?
pReader
->
pLastFD
:
pReader
->
pDataFD
;
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
iSubBlock
];
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
iSubBlock
];
SArray
*
aBlockCol
=
NULL
;
int32_t
code
=
0
;
int32_t
code
=
0
;
int64_t
offset
;
int64_t
offset
;
int64_t
size
;
int64_t
size
;
...
@@ -766,9 +837,15 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
...
@@ -766,9 +837,15 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
tBlockDataReset
(
pBlockData
);
tBlockDataReset
(
pBlockData
);
pBlockData
->
nRow
=
pSubBlock
->
nRow
;
pBlockData
->
nRow
=
pSubBlock
->
nRow
;
// TSDBKEY
// TSDBKEY and SBlockCol
offset
=
pSubBlock
->
offset
+
sizeof
(
SBlockDataHdr
);
if
(
nCol
==
1
)
{
offset
=
pSubBlock
->
offset
+
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
);
size
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
size
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
}
else
{
offset
=
pSubBlock
->
offset
;
size
=
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
)
+
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
}
code
=
tsdbRealloc
(
ppBuf1
,
size
);
code
=
tsdbRealloc
(
ppBuf1
,
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
@@ -787,31 +864,47 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
...
@@ -787,31 +864,47 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
goto
_err
;
goto
_err
;
}
}
code
=
tsdbRecoverBlockDataKey
(
pBlockData
,
pSubBlock
,
*
ppBuf1
,
ppBuf2
);
if
(
nCol
==
1
)
{
code
=
tsdbReadBlockDataKey
(
pBlockData
,
pSubBlock
,
*
ppBuf1
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
// OTHER
goto
_exit
;
SBlockCol
blockCol
;
}
else
{
SBlockCol
*
pBlockCol
=
&
blockCol
;
aBlockCol
=
taosArrayInit
(
0
,
sizeof
(
SBlockCol
));
SColData
*
pColData
;
if
(
aBlockCol
==
NULL
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
code
=
tsdbReadBlockCol
(
pSubBlock
,
*
ppBuf1
,
aBlockCol
);
if
(
code
)
goto
_err
;
code
=
tsdbReadBlockDataKey
(
pBlockData
,
pSubBlock
,
*
ppBuf1
+
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
),
ppBuf2
);
if
(
code
)
goto
_err
;
}
for
(
int32_t
iCol
=
1
;
iCol
<
nCol
;
iCol
++
)
{
for
(
int32_t
iCol
=
1
;
iCol
<
nCol
;
iCol
++
)
{
int16_t
cid
=
aColId
[
iCol
];
void
*
p
=
taosArraySearch
(
aBlockCol
,
&
(
SBlockCol
){.
cid
=
aColId
[
iCol
]},
tBlockColCmprFn
,
TD_EQ
);
if
(
p
)
{
SBlockCol
*
pBlockCol
=
(
SBlockCol
*
)
p
;
SColData
*
pColData
;
ASSERT
(
pBlockCol
->
flag
&&
pBlockCol
->
flag
!=
HAS_NONE
);
if
(
tMapDataSearch
(
&
pSubBlock
->
mBlockCol
,
&
(
SBlockCol
){.
cid
=
cid
},
tGetBlockCol
,
tBlockColCmprFn
,
pBlockCol
)
==
0
)
{
code
=
tBlockDataAddColData
(
pBlockData
,
taosArrayGetSize
(
pBlockData
->
aColDataP
),
&
pColData
);
code
=
tBlockDataAddColData
(
pBlockData
,
taosArrayGetSize
(
pBlockData
->
aColDataP
),
&
pColData
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
tColDataReset
(
pColData
,
pBlockCol
->
cid
,
pBlockCol
->
type
);
tColDataReset
(
pColData
,
pBlockCol
->
cid
,
pBlockCol
->
type
,
0
);
if
(
pBlockCol
->
flag
==
HAS_NULL
)
{
if
(
pBlockCol
->
flag
==
HAS_NULL
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pSubBlock
->
nRow
;
iRow
++
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pSubBlock
->
nRow
;
iRow
++
)
{
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NULL
(
pBlockCol
->
cid
,
pBlockCol
->
type
));
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NULL
(
pBlockCol
->
cid
,
pBlockCol
->
type
));
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
}
}
else
{
}
else
{
offset
=
pSubBlock
->
offset
+
sizeof
(
SBlockDataHdr
)
+
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
offset
=
pSubBlock
->
offset
+
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
)
+
pSubBlock
->
szVersion
+
sizeof
(
TSCKSUM
)
+
pBlockCol
->
offset
;
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
)
+
pBlockCol
->
offset
;
size
=
pBlockCol
->
bsize
+
pBlockCol
->
csiz
e
+
sizeof
(
TSCKSUM
);
size
=
pBlockCol
->
szBitmap
+
pBlockCol
->
szOffset
+
pBlockCol
->
szValu
e
+
sizeof
(
TSCKSUM
);
code
=
tsdbRealloc
(
ppBuf1
,
size
);
code
=
tsdbRealloc
(
ppBuf1
,
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
@@ -833,14 +926,18 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
...
@@ -833,14 +926,18 @@ static int32_t tsdbReadColDataImpl(SDataFReader *pReader, SBlockIdx *pBlockIdx,
goto
_err
;
goto
_err
;
}
}
code
=
tsdbRe
coverColData
(
pBlockData
,
pSubBlock
,
pBlockCol
,
pColData
,
*
ppBuf1
,
ppBuf2
);
code
=
tsdbRe
adColDataImpl
(
pSubBlock
,
pBlockCol
,
pColData
,
*
ppBuf1
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
}
}
}
}
}
_exit:
taosArrayDestroy
(
aBlockCol
);
return
code
;
return
code
;
_err:
_err:
taosArrayDestroy
(
aBlockCol
);
return
code
;
return
code
;
}
}
...
@@ -855,7 +952,7 @@ int32_t tsdbReadColData(SDataFReader *pReader, SBlockIdx *pBlockIdx, SBlock *pBl
...
@@ -855,7 +952,7 @@ int32_t tsdbReadColData(SDataFReader *pReader, SBlockIdx *pBlockIdx, SBlock *pBl
if
(
!
ppBuf1
)
ppBuf1
=
&
pBuf1
;
if
(
!
ppBuf1
)
ppBuf1
=
&
pBuf1
;
if
(
!
ppBuf2
)
ppBuf2
=
&
pBuf2
;
if
(
!
ppBuf2
)
ppBuf2
=
&
pBuf2
;
code
=
tsdbRead
ColDataImpl
(
pReader
,
pBlockIdx
,
pBlock
,
0
,
aColId
,
nCol
,
pBlockData
,
ppBuf1
,
ppBuf2
);
code
=
tsdbRead
SubColData
(
pReader
,
pBlockIdx
,
pBlock
,
0
,
aColId
,
nCol
,
pBlockData
,
ppBuf1
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
if
(
pBlock
->
nSubBlock
>
1
)
{
if
(
pBlock
->
nSubBlock
>
1
)
{
...
@@ -863,7 +960,7 @@ int32_t tsdbReadColData(SDataFReader *pReader, SBlockIdx *pBlockIdx, SBlock *pBl
...
@@ -863,7 +960,7 @@ int32_t tsdbReadColData(SDataFReader *pReader, SBlockIdx *pBlockIdx, SBlock *pBl
SBlockData
*
pBlockData2
=
&
(
SBlockData
){
0
};
SBlockData
*
pBlockData2
=
&
(
SBlockData
){
0
};
for
(
int32_t
iSubBlock
=
1
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
for
(
int32_t
iSubBlock
=
1
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
code
=
tsdbRead
ColDataImpl
(
pReader
,
pBlockIdx
,
pBlock
,
iSubBlock
,
aColId
,
nCol
,
pBlockData1
,
ppBuf1
,
ppBuf2
);
code
=
tsdbRead
SubColData
(
pReader
,
pBlockIdx
,
pBlock
,
iSubBlock
,
aColId
,
nCol
,
pBlockData1
,
ppBuf1
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
code
=
tBlockDataCopy
(
pBlockData
,
pBlockData2
);
code
=
tBlockDataCopy
(
pBlockData
,
pBlockData2
);
...
@@ -904,7 +1001,7 @@ static int32_t tsdbReadSubBlockData(SDataFReader *pReader, SBlockIdx *pBlockIdx,
...
@@ -904,7 +1001,7 @@ static int32_t tsdbReadSubBlockData(SDataFReader *pReader, SBlockIdx *pBlockIdx,
int64_t
n
;
int64_t
n
;
TdFilePtr
pFD
=
pBlock
->
last
?
pReader
->
pLastFD
:
pReader
->
pDataFD
;
TdFilePtr
pFD
=
pBlock
->
last
?
pReader
->
pLastFD
:
pReader
->
pDataFD
;
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
iSubBlock
];
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
iSubBlock
];
S
BlockCol
*
pBlockCol
=
&
(
SBlockCol
){
0
}
;
S
Array
*
aBlockCol
=
NULL
;
tBlockDataReset
(
pBlockData
);
tBlockDataReset
(
pBlockData
);
...
@@ -929,41 +1026,52 @@ static int32_t tsdbReadSubBlockData(SDataFReader *pReader, SBlockIdx *pBlockIdx,
...
@@ -929,41 +1026,52 @@ static int32_t tsdbReadSubBlockData(SDataFReader *pReader, SBlockIdx *pBlockIdx,
goto
_err
;
goto
_err
;
}
}
// recover
pBlockData
->
nRow
=
pSubBlock
->
nRow
;
pBlockData
->
nRow
=
pSubBlock
->
nRow
;
p
=
*
ppBuf1
+
sizeof
(
SBlockDataHdr
);
code
=
tsdbRecoverBlockDataKey
(
pBlockData
,
pSubBlock
,
p
,
ppBuf2
);
// TSDBKEY
p
=
*
ppBuf1
+
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
);
code
=
tsdbReadBlockDataKey
(
pBlockData
,
pSubBlock
,
p
,
ppBuf2
);
if
(
code
)
goto
_err
;
// COLUMNS
aBlockCol
=
taosArrayInit
(
0
,
sizeof
(
SBlockCol
));
if
(
aBlockCol
==
NULL
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
code
=
tsdbReadBlockCol
(
pSubBlock
,
*
ppBuf1
,
aBlockCol
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
p
=
p
+
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
for
(
int32_t
iBlockCol
=
0
;
iBlockCol
<
pSubBlock
->
mBlockCol
.
nItem
;
iBlockCol
++
)
{
for
(
int32_t
iBlockCol
=
0
;
iBlockCol
<
taosArrayGetSize
(
aBlockCol
)
;
iBlockCol
++
)
{
SColData
*
pColData
;
SColData
*
pColData
;
SBlockCol
*
pBlockCol
=
(
SBlockCol
*
)
taosArrayGet
(
aBlockCol
,
iBlockCol
);
tMapDataGetItemByIdx
(
&
pSubBlock
->
mBlockCol
,
iBlockCol
,
pBlockCol
,
tGetBlockCol
);
ASSERT
(
pBlockCol
->
flag
&&
pBlockCol
->
flag
!=
HAS_NONE
);
ASSERT
(
pBlockCol
->
flag
&&
pBlockCol
->
flag
!=
HAS_NONE
);
code
=
tBlockDataAddColData
(
pBlockData
,
iBlockCol
,
&
pColData
);
code
=
tBlockDataAddColData
(
pBlockData
,
iBlockCol
,
&
pColData
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
tColDataReset
(
pColData
,
pBlockCol
->
cid
,
pBlockCol
->
type
);
tColDataReset
(
pColData
,
pBlockCol
->
cid
,
pBlockCol
->
type
,
0
);
if
(
pBlockCol
->
flag
==
HAS_NULL
)
{
if
(
pBlockCol
->
flag
==
HAS_NULL
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pSubBlock
->
nRow
;
iRow
++
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pSubBlock
->
nRow
;
iRow
++
)
{
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NULL
(
pBlockCol
->
cid
,
pBlockCol
->
type
));
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NULL
(
pBlockCol
->
cid
,
pBlockCol
->
type
));
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
}
}
else
{
}
else
{
code
=
tsdbRecoverColData
(
pBlockData
,
pSubBlock
,
pBlockCol
,
pColData
,
p
,
ppBuf2
);
p
=
*
ppBuf1
+
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
)
+
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
)
+
pBlockCol
->
offset
;
code
=
tsdbReadColDataImpl
(
pSubBlock
,
pBlockCol
,
pColData
,
p
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
p
=
p
+
pBlockCol
->
bsize
+
pBlockCol
->
csize
+
sizeof
(
TSCKSUM
);
}
}
}
}
taosArrayDestroy
(
aBlockCol
);
return
code
;
return
code
;
_err:
_err:
tsdbError
(
"vgId:%d tsdb read sub block data failed since %s"
,
TD_VID
(
pReader
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbError
(
"vgId:%d tsdb read sub block data failed since %s"
,
TD_VID
(
pReader
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
taosArrayDestroy
(
aBlockCol
);
return
code
;
return
code
;
}
}
...
@@ -1189,6 +1297,8 @@ int32_t tsdbDataFWriterClose(SDataFWriter **ppWriter, int8_t sync) {
...
@@ -1189,6 +1297,8 @@ int32_t tsdbDataFWriterClose(SDataFWriter **ppWriter, int8_t sync) {
int32_t
code
=
0
;
int32_t
code
=
0
;
STsdb
*
pTsdb
=
(
*
ppWriter
)
->
pTsdb
;
STsdb
*
pTsdb
=
(
*
ppWriter
)
->
pTsdb
;
if
(
*
ppWriter
==
NULL
)
goto
_exit
;
if
(
sync
)
{
if
(
sync
)
{
if
(
taosFsyncFile
((
*
ppWriter
)
->
pHeadFD
)
<
0
)
{
if
(
taosFsyncFile
((
*
ppWriter
)
->
pHeadFD
)
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
...
@@ -1232,6 +1342,7 @@ int32_t tsdbDataFWriterClose(SDataFWriter **ppWriter, int8_t sync) {
...
@@ -1232,6 +1342,7 @@ int32_t tsdbDataFWriterClose(SDataFWriter **ppWriter, int8_t sync) {
}
}
taosMemoryFree
(
*
ppWriter
);
taosMemoryFree
(
*
ppWriter
);
_exit:
*
ppWriter
=
NULL
;
*
ppWriter
=
NULL
;
return
code
;
return
code
;
...
@@ -1366,246 +1477,390 @@ _err:
...
@@ -1366,246 +1477,390 @@ _err:
return
code
;
return
code
;
}
}
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
uint8_t
**
ppBuf2
,
static
void
tsdbUpdateBlockInfo
(
SBlockData
*
pBlockData
,
SBlock
*
pBlock
)
{
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int8_t
cmprAlg
)
{
int32_t
code
=
0
;
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
pBlock
->
nSubBlock
++
];
SBlockCol
*
pBlockCol
=
&
(
SBlockCol
){
0
};
int64_t
size
;
int64_t
n
;
TdFilePtr
pFileFD
=
pBlock
->
last
?
pWriter
->
pLastFD
:
pWriter
->
pDataFD
;
SBlockDataHdr
hdr
=
{.
delimiter
=
TSDB_FILE_DLMT
,
.
suid
=
pBlockIdx
->
suid
,
.
uid
=
pBlockIdx
->
uid
};
TSCKSUM
cksm
;
uint8_t
*
p
;
int64_t
offset
;
uint8_t
*
pBuf1
=
NULL
;
uint8_t
*
pBuf2
=
NULL
;
if
(
!
ppBuf1
)
ppBuf1
=
&
pBuf1
;
if
(
!
ppBuf2
)
ppBuf2
=
&
pBuf2
;
TSKEY
lastKey
=
TSKEY_MIN
;
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
TSDBKEY
key
=
TSDBROW_KEY
(
&
tsdbRowFromBlockData
(
pBlockData
,
iRow
));
TSDBKEY
key
=
{.
ts
=
pBlockData
->
aTSKEY
[
iRow
],
.
version
=
pBlockData
->
aVersion
[
iRow
]};
if
(
iRow
==
0
)
{
if
(
iRow
==
0
)
{
pBlock
->
minKey
=
MIN_TSDBKEY
(
pBlock
->
minKey
,
key
);
if
(
tsdbKeyCmprFn
(
&
pBlock
->
minKey
,
&
key
)
>
0
)
{
pBlock
->
minKey
=
key
;
}
}
else
{
if
(
pBlockData
->
aTSKEY
[
iRow
]
==
pBlockData
->
aTSKEY
[
iRow
-
1
])
{
pBlock
->
hasDup
=
1
;
}
}
}
if
(
iRow
==
pBlockData
->
nRow
-
1
)
{
if
(
iRow
==
pBlockData
->
nRow
-
1
&&
tsdbKeyCmprFn
(
&
pBlock
->
maxKey
,
&
key
)
<
0
)
{
pBlock
->
maxKey
=
MAX_TSDBKEY
(
pBlock
->
maxKey
,
key
)
;
pBlock
->
maxKey
=
key
;
}
}
pBlock
->
minVersion
=
TMIN
(
pBlock
->
minVersion
,
key
.
version
);
pBlock
->
minVersion
=
TMIN
(
pBlock
->
minVersion
,
key
.
version
);
pBlock
->
maxVersion
=
TMAX
(
pBlock
->
maxVersion
,
key
.
version
);
pBlock
->
maxVersion
=
TMAX
(
pBlock
->
maxVersion
,
key
.
version
);
if
(
key
.
ts
==
lastKey
)
{
pBlock
->
hasDup
=
1
;
}
lastKey
=
key
.
ts
;
}
}
pBlock
->
nRow
+=
pBlockData
->
nRow
;
pBlock
->
nRow
+=
pBlockData
->
nRow
;
}
pSubBlock
->
nRow
=
pBlockData
->
nRow
;
static
int32_t
tsdbWriteBlockDataKey
(
SSubBlock
*
pSubBlock
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
int64_t
*
nDataP
,
pSubBlock
->
cmprAlg
=
cmprAlg
;
uint8_t
**
ppBuf2
)
{
if
(
pBlock
->
last
)
{
int32_t
code
=
0
;
pSubBlock
->
offset
=
pWriter
->
wSet
.
fLast
.
size
;
int64_t
size
;
int64_t
tsize
;
if
(
pSubBlock
->
cmprAlg
==
NO_COMPRESSION
)
{
pSubBlock
->
szVersion
=
sizeof
(
int64_t
)
*
pSubBlock
->
nRow
;
pSubBlock
->
szTSKEY
=
sizeof
(
TSKEY
)
*
pSubBlock
->
nRow
;
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
));
if
(
code
)
goto
_err
;
// VERSION
memcpy
(
*
ppBuf1
+
*
nDataP
,
pBlockData
->
aVersion
,
pSubBlock
->
szVersion
);
// TSKEY
memcpy
(
*
ppBuf1
+
*
nDataP
+
pSubBlock
->
szVersion
,
pBlockData
->
aTSKEY
,
pSubBlock
->
szTSKEY
);
}
else
{
}
else
{
pSubBlock
->
offset
=
pWriter
->
wSet
.
fData
.
size
;
size
=
(
sizeof
(
int64_t
)
+
sizeof
(
TSKEY
))
*
pSubBlock
->
nRow
+
COMP_OVERFLOW_BYTES
*
2
;
}
pSubBlock
->
szBlock
=
0
;
// HDR
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
size
+
sizeof
(
TSCKSUM
));
n
=
taosWriteFile
(
pFileFD
,
&
hdr
,
sizeof
(
hdr
));
if
(
code
)
goto
_err
;
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
}
pSubBlock
->
szBlock
+=
n
;
// TSDBKEY
tsize
=
sizeof
(
int64_t
)
*
pSubBlock
->
nRow
+
COMP_OVERFLOW_BYTES
;
if
(
cmprAlg
==
NO_COMPRESSION
)
{
if
(
pSubBlock
->
cmprAlg
==
TWO_STAGE_COMP
)
{
cksm
=
0
;
code
=
tsdbRealloc
(
ppBuf2
,
tsize
);
if
(
code
)
goto
_err
;
}
// version
// VERSION
pSubBlock
->
szVersion
=
sizeof
(
int64_t
)
*
pBlockData
->
nRow
;
pSubBlock
->
szVersion
=
n
=
taosWriteFile
(
pFileFD
,
pBlockData
->
aVersion
,
pSubBlock
->
szVersion
);
tsCompressBigint
((
char
*
)
pBlockData
->
aVersion
,
sizeof
(
int64_t
)
*
pBlockData
->
nRow
,
pBlockData
->
nRow
,
if
(
n
<
0
)
{
*
ppBuf1
+
*
nDataP
,
size
,
pSubBlock
->
cmprAlg
,
*
ppBuf2
,
tsize
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
if
(
pSubBlock
->
szVersion
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
}
}
cksm
=
taosCalcChecksum
(
cksm
,
(
uint8_t
*
)
pBlockData
->
aVersion
,
pSubBlock
->
szVersion
);
// TSKEY
// TSKEY
pSubBlock
->
szTSKEY
=
sizeof
(
TSKEY
)
*
pBlockData
->
nRow
;
pSubBlock
->
szTSKEY
=
tsCompressTimestamp
((
char
*
)
pBlockData
->
aTSKEY
,
sizeof
(
TSKEY
)
*
pBlockData
->
nRow
,
n
=
taosWriteFile
(
pFileFD
,
pBlockData
->
aTSKEY
,
pSubBlock
->
szTSKEY
);
pBlockData
->
nRow
,
*
ppBuf1
+
*
nDataP
+
pSubBlock
->
szVersion
,
if
(
n
<
0
)
{
size
-
pSubBlock
->
szVersion
,
pSubBlock
->
cmprAlg
,
*
ppBuf2
,
tsize
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
if
(
pSubBlock
->
szTSKEY
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
}
}
cksm
=
taosCalcChecksum
(
cksm
,
(
uint8_t
*
)
pBlockData
->
aTSKEY
,
pSubBlock
->
szTSKEY
);
// cksm
ASSERT
(
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
<=
size
);
size
=
sizeof
(
cksm
);
n
=
taosWriteFile
(
pFileFD
,
(
uint8_t
*
)
&
cksm
,
size
);
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
}
}
}
else
{
ASSERT
(
cmprAlg
==
ONE_STAGE_COMP
||
cmprAlg
==
TWO_STAGE_COMP
);
size
=
(
sizeof
(
int64_t
)
+
sizeof
(
TSKEY
))
*
pBlockData
->
nRow
+
COMP_OVERFLOW_BYTES
*
2
+
sizeof
(
TSCKSUM
);
// checksum
size
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
taosCalcChecksumAppend
(
0
,
*
ppBuf1
+
*
nDataP
,
size
);
*
nDataP
+=
size
;
return
code
;
code
=
tsdbRealloc
(
ppBuf1
,
size
);
_err:
return
code
;
}
static
int32_t
tsdbWriteColData
(
SColData
*
pColData
,
SBlockCol
*
pBlockCol
,
SSubBlock
*
pSubBlock
,
uint8_t
**
ppBuf1
,
int64_t
*
nDataP
,
uint8_t
**
ppBuf2
)
{
int32_t
code
=
0
;
int64_t
size
;
int64_t
n
=
0
;
// BITMAP
if
(
pColData
->
flag
!=
HAS_VALUE
)
{
size
=
BIT2_SIZE
(
pColData
->
nVal
)
+
COMP_OVERFLOW_BYTES
;
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
n
+
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
if
(
cmprAlg
==
TWO_STAGE_COMP
)
{
code
=
tsdbRealloc
(
ppBuf2
,
size
);
code
=
tsdbRealloc
(
ppBuf2
,
size
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
}
// version
pBlockCol
->
szBitmap
=
n
=
tsCompressBigint
((
char
*
)
pBlockData
->
aVersion
,
sizeof
(
int64_t
)
*
pBlockData
->
nRow
,
pBlockData
->
nRow
,
*
ppBuf1
,
tsCompressTinyint
((
char
*
)
pColData
->
pBitMap
,
BIT2_SIZE
(
pColData
->
nVal
),
BIT2_SIZE
(
pColData
->
nVal
)
,
size
,
cmprAlg
,
*
ppBuf2
,
size
);
*
ppBuf1
+
*
nDataP
+
n
,
size
,
TWO_STAGE_COMP
,
*
ppBuf2
,
size
);
if
(
n
<=
0
)
{
if
(
pBlockCol
->
szBitmap
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
}
}
pSubBlock
->
szVersion
=
n
;
}
else
{
pBlockCol
->
szBitmap
=
0
;
}
n
+=
pBlockCol
->
szBitmap
;
// TSKEY
// OFFSET
n
=
tsCompressTimestamp
((
char
*
)
pBlockData
->
aTSKEY
,
sizeof
(
TSKEY
)
*
pBlockData
->
nRow
,
pBlockData
->
nRow
,
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
{
*
ppBuf1
+
pSubBlock
->
szVersion
,
size
-
pSubBlock
->
szVersion
,
cmprAlg
,
*
ppBuf2
,
size
);
size
=
sizeof
(
int32_t
)
*
pColData
->
nVal
+
COMP_OVERFLOW_BYTES
;
if
(
n
<=
0
)
{
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
n
+
size
);
if
(
code
)
goto
_err
;
code
=
tsdbRealloc
(
ppBuf2
,
size
);
if
(
code
)
goto
_err
;
pBlockCol
->
szOffset
=
tsCompressInt
((
char
*
)
pColData
->
aOffset
,
sizeof
(
int32_t
)
*
pColData
->
nVal
,
pColData
->
nVal
,
*
ppBuf1
+
*
nDataP
+
n
,
size
,
TWO_STAGE_COMP
,
*
ppBuf2
,
size
);
if
(
pBlockCol
->
szOffset
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
}
}
pSubBlock
->
szTSKEY
=
n
;
}
else
{
pBlockCol
->
szOffset
=
0
;
}
n
+=
pBlockCol
->
szOffset
;
// VALUE
if
(
pSubBlock
->
cmprAlg
==
NO_COMPRESSION
)
{
pBlockCol
->
szValue
=
pColData
->
nData
;
// cksm
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
n
+
pBlockCol
->
szValue
+
sizeof
(
TSCKSUM
));
n
=
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
);
if
(
code
)
goto
_err
;
ASSERT
(
n
<=
size
);
taosCalcChecksumAppend
(
0
,
*
ppBuf1
,
n
);
// write
memcpy
(
*
ppBuf1
+
*
nDataP
+
n
,
pColData
->
pData
,
pBlockCol
->
szValue
);
n
=
taosWriteFile
(
pFileFD
,
*
ppBuf1
,
n
);
}
else
{
if
(
n
<
0
)
{
size
=
pColData
->
nData
+
COMP_OVERFLOW_BYTES
;
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
tsdbRealloc
(
ppBuf1
,
*
nDataP
+
n
+
size
+
sizeof
(
TSCKSUM
));
if
(
code
)
goto
_err
;
if
(
pSubBlock
->
cmprAlg
==
TWO_STAGE_COMP
)
{
code
=
tsdbRealloc
(
ppBuf2
,
size
);
if
(
code
)
goto
_err
;
}
pBlockCol
->
szValue
=
tDataTypes
[
pColData
->
type
].
compFunc
((
char
*
)
pColData
->
pData
,
pColData
->
nData
,
pColData
->
nVal
,
*
ppBuf1
+
*
nDataP
+
n
,
size
,
pSubBlock
->
cmprAlg
,
*
ppBuf2
,
size
);
if
(
pBlockCol
->
szValue
<=
0
)
{
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
goto
_err
;
}
}
}
}
pSubBlock
->
szBlock
+=
(
pSubBlock
->
szVersion
+
pSubBlock
->
szTSKEY
+
sizeof
(
TSCKSUM
));
n
+=
pBlockCol
->
szValue
;
pBlockCol
->
szOrigin
=
pColData
->
nData
;
// other columns
// checksum
offset
=
0
;
n
+=
sizeof
(
TSCKSUM
);
tMapDataReset
(
&
pSubBlock
->
mBlockCol
);
taosCalcChecksumAppend
(
0
,
*
ppBuf1
+
*
nDataP
,
n
);
for
(
int32_t
iCol
=
0
;
iCol
<
taosArrayGetSize
(
pBlockData
->
aColDataP
);
iCol
++
)
{
SColData
*
pColData
=
(
SColData
*
)
taosArrayGetP
(
pBlockData
->
aColDataP
,
iCol
);
ASSERT
(
pColData
->
flag
)
;
*
nDataP
+=
n
;
if
(
pColData
->
flag
==
HAS_NONE
)
continu
e
;
return
cod
e
;
pBlockCol
->
cid
=
pColData
->
cid
;
_err:
pBlockCol
->
type
=
pColData
->
typ
e
;
return
cod
e
;
pBlockCol
->
flag
=
pColData
->
flag
;
}
if
(
pColData
->
flag
!=
HAS_NULL
)
{
static
int32_t
tsdbWriteBlockDataImpl
(
TdFilePtr
pFD
,
SSubBlock
*
pSubBlock
,
SBlockDataHdr
hdr
,
SArray
*
aBlockCol
,
cksm
=
0
;
uint8_t
*
pData
,
int64_t
nData
,
uint8_t
**
ppBuf
)
{
pBlockCol
->
offset
=
offset
;
int32_t
code
=
0
;
int32_t
nBlockCol
=
taosArrayGetSize
(
aBlockCol
);
int64_t
size
;
int64_t
n
;
// bitmap
// HDR + SArray<SBlockCol>
if
(
pColData
->
flag
==
HAS_VALUE
)
{
pSubBlock
->
szBlockCol
=
sizeof
(
hdr
);
pBlockCol
->
bsize
=
0
;
for
(
int32_t
iBlockCol
=
0
;
iBlockCol
<
nBlockCol
;
iBlockCol
++
)
{
}
else
{
pSubBlock
->
szBlockCol
+=
tPutBlockCol
(
NULL
,
taosArrayGet
(
aBlockCol
,
iBlockCol
));
pBlockCol
->
bsize
=
BIT2_SIZE
(
pBlockData
->
nRow
);
n
=
taosWriteFile
(
pFileFD
,
pColData
->
pBitMap
,
pBlockCol
->
bsize
);
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
}
}
cksm
=
taosCalcChecksum
(
cksm
,
pColData
->
pBitMap
,
n
);
code
=
tsdbRealloc
(
ppBuf
,
pSubBlock
->
szBlock
+
sizeof
(
TSCKSUM
));
if
(
code
)
goto
_err
;
n
=
0
;
memcpy
(
*
ppBuf
,
&
hdr
,
sizeof
(
hdr
));
n
+=
sizeof
(
hdr
);
for
(
int32_t
iBlockCol
=
0
;
iBlockCol
<
nBlockCol
;
iBlockCol
++
)
{
n
+=
tPutBlockCol
(
*
ppBuf
+
n
,
taosArrayGet
(
aBlockCol
,
iBlockCol
));
}
}
taosCalcChecksumAppend
(
0
,
*
ppBuf
,
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
));
// data
ASSERT
(
n
==
pSubBlock
->
szBlockCol
);
if
(
cmprAlg
==
NO_COMPRESSION
)
{
// data
n
=
taosWriteFile
(
pFD
,
*
ppBuf
,
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
));
n
=
taosWriteFile
(
pFileFD
,
pColData
->
pData
,
pColData
->
nData
);
if
(
n
<
0
)
{
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
goto
_err
;
}
}
pBlockCol
->
csize
=
n
;
pBlockCol
->
osize
=
n
;
// checksum
// SBlockData
cksm
=
taosCalcChecksum
(
cksm
,
pColData
->
pData
,
pColData
->
nData
);
n
=
taosWriteFile
(
pFD
,
pData
,
nData
);
n
=
taosWriteFile
(
pFileFD
,
&
cksm
,
sizeof
(
cksm
));
if
(
n
<
0
)
{
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
goto
_err
;
}
}
}
else
{
size
=
pColData
->
nData
+
COMP_OVERFLOW_BYTES
+
sizeof
(
TSCKSUM
);
code
=
tsdbRealloc
(
ppBuf1
,
size
);
return
code
;
if
(
code
)
goto
_err
;
if
(
cmprAlg
==
TWO_STAGE_COMP
)
{
_err:
code
=
tsdbRealloc
(
ppBuf2
,
size
);
return
code
;
if
(
code
)
goto
_err
;
}
static
void
tsdbCalcColDataSMA
(
SColData
*
pColData
,
SColumnDataAgg
*
pColAgg
)
{
SColVal
colVal
;
SColVal
*
pColVal
=
&
colVal
;
*
pColAgg
=
(
SColumnDataAgg
){.
colId
=
pColData
->
cid
};
for
(
int32_t
iVal
=
0
;
iVal
<
pColData
->
nVal
;
iVal
++
)
{
tColDataGetValue
(
pColData
,
iVal
,
pColVal
);
if
(
pColVal
->
isNone
||
pColVal
->
isNull
)
{
pColAgg
->
numOfNull
++
;
}
else
{
// TODO:
ASSERT
(
0
);
}
}
}
}
// data
static
int32_t
tsdbWriteBlockSMA
(
TdFilePtr
pFD
,
SBlockData
*
pBlockData
,
SSubBlock
*
pSubBlock
,
uint8_t
**
ppBuf
)
{
n
=
tDataTypes
[
pColData
->
type
].
compFunc
(
pColData
->
pData
,
pColData
->
nData
,
pBlockData
->
nRow
,
*
ppBuf1
,
size
,
int32_t
code
=
0
;
cmprAlg
,
*
ppBuf2
,
size
);
int64_t
n
;
if
(
n
<=
0
)
{
SColData
*
pColData
;
code
=
TSDB_CODE_COMPRESS_ERROR
;
goto
_err
;
// prepare
pSubBlock
->
nSma
=
0
;
for
(
int32_t
iColData
=
0
;
iColData
<
taosArrayGetSize
(
pBlockData
->
aColDataP
);
iColData
++
)
{
pColData
=
(
SColData
*
)
taosArrayGetP
(
pBlockData
->
aColDataP
,
iColData
);
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
)
||
(
!
pColData
->
smaOn
))
continue
;
pSubBlock
->
nSma
++
;
}
}
pBlockCol
->
csize
=
n
;
if
(
pSubBlock
->
nSma
==
0
)
goto
_exit
;
pBlockCol
->
osize
=
pColData
->
nData
;
// cksm
// calc
n
+=
sizeof
(
TSCKSUM
);
code
=
tsdbRealloc
(
ppBuf
,
sizeof
(
SColumnDataAgg
)
*
pSubBlock
->
nSma
+
sizeof
(
TSCKSUM
));
ASSERT
(
n
<=
size
);
if
(
code
)
goto
_err
;
taosCalcChecksumAppend
(
cksm
,
*
ppBuf1
,
n
);
n
=
0
;
for
(
int32_t
iColData
=
0
;
iColData
<
taosArrayGetSize
(
pBlockData
->
aColDataP
);
iColData
++
)
{
pColData
=
(
SColData
*
)
taosArrayGetP
(
pBlockData
->
aColDataP
,
iColData
);
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
)
||
(
!
pColData
->
smaOn
))
continue
;
tsdbCalcColDataSMA
(
pColData
,
&
((
SColumnDataAgg
*
)(
*
ppBuf
))[
n
]);
n
++
;
}
taosCalcChecksumAppend
(
0
,
*
ppBuf
,
sizeof
(
SColumnDataAgg
)
*
pSubBlock
->
nSma
+
sizeof
(
TSCKSUM
));
// write
// write
n
=
taosWriteFile
(
pFileFD
,
*
ppBuf1
,
n
);
n
=
taosWriteFile
(
pFD
,
*
ppBuf
,
sizeof
(
SColumnDataAgg
)
*
pSubBlock
->
nSma
+
sizeof
(
TSCKSUM
)
);
if
(
n
<
0
)
{
if
(
n
<
0
)
{
code
=
TAOS_SYSTEM_ERROR
(
errno
);
code
=
TAOS_SYSTEM_ERROR
(
errno
);
goto
_err
;
goto
_err
;
}
}
_exit:
return
code
;
_err:
return
code
;
}
int32_t
tsdbWriteBlockData
(
SDataFWriter
*
pWriter
,
SBlockData
*
pBlockData
,
uint8_t
**
ppBuf1
,
uint8_t
**
ppBuf2
,
SBlockIdx
*
pBlockIdx
,
SBlock
*
pBlock
,
int8_t
cmprAlg
)
{
int32_t
code
=
0
;
SSubBlock
*
pSubBlock
=
&
pBlock
->
aSubBlock
[
pBlock
->
nSubBlock
++
];
SBlockCol
blockCol
;
SBlockCol
*
pBlockCol
=
&
blockCol
;
int64_t
n
;
TdFilePtr
pFileFD
=
pBlock
->
last
?
pWriter
->
pLastFD
:
pWriter
->
pDataFD
;
SBlockDataHdr
hdr
=
{.
delimiter
=
TSDB_FILE_DLMT
,
.
suid
=
pBlockIdx
->
suid
,
.
uid
=
pBlockIdx
->
uid
};
uint8_t
*
p
;
int64_t
nData
;
uint8_t
*
pBuf1
=
NULL
;
uint8_t
*
pBuf2
=
NULL
;
SArray
*
aBlockCol
=
NULL
;
if
(
!
ppBuf1
)
ppBuf1
=
&
pBuf1
;
if
(
!
ppBuf2
)
ppBuf2
=
&
pBuf2
;
tsdbUpdateBlockInfo
(
pBlockData
,
pBlock
);
pSubBlock
->
nRow
=
pBlockData
->
nRow
;
pSubBlock
->
cmprAlg
=
cmprAlg
;
if
(
pBlock
->
last
)
{
pSubBlock
->
offset
=
pWriter
->
wSet
.
fLast
.
size
;
}
else
{
pSubBlock
->
offset
=
pWriter
->
wSet
.
fData
.
size
;
}
}
// state
// ======================= BLOCK DATA =======================
offset
=
offset
+
pBlockCol
->
bsize
+
pBlockCol
->
csize
+
sizeof
(
TSCKSUM
);
// TSDBKEY
pSubBlock
->
szBlock
=
pSubBlock
->
szBlock
+
pBlockCol
->
bsize
+
pBlockCol
->
csize
+
sizeof
(
TSCKSUM
);
nData
=
0
;
code
=
tsdbWriteBlockDataKey
(
pSubBlock
,
pBlockData
,
ppBuf1
,
&
nData
,
ppBuf2
);
if
(
code
)
goto
_err
;
// COLUMNS
aBlockCol
=
taosArrayInit
(
taosArrayGetSize
(
pBlockData
->
aColDataP
),
sizeof
(
SBlockCol
));
if
(
aBlockCol
==
NULL
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
}
int32_t
offset
=
0
;
for
(
int32_t
iCol
=
0
;
iCol
<
taosArrayGetSize
(
pBlockData
->
aColDataP
);
iCol
++
)
{
SColData
*
pColData
=
(
SColData
*
)
taosArrayGetP
(
pBlockData
->
aColDataP
,
iCol
);
code
=
tMapDataPutItem
(
&
pSubBlock
->
mBlockCol
,
pBlockCol
,
tPutBlockCol
);
ASSERT
(
pColData
->
flag
);
if
(
pColData
->
flag
==
HAS_NONE
)
continue
;
pBlockCol
->
cid
=
pColData
->
cid
;
pBlockCol
->
type
=
pColData
->
type
;
pBlockCol
->
flag
=
pColData
->
flag
;
if
(
pColData
->
flag
!=
HAS_NULL
)
{
code
=
tsdbWriteColData
(
pColData
,
pBlockCol
,
pSubBlock
,
ppBuf1
,
&
nData
,
ppBuf2
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
pBlockCol
->
offset
=
offset
;
offset
=
offset
+
pBlockCol
->
szBitmap
+
pBlockCol
->
szOffset
+
pBlockCol
->
szValue
+
sizeof
(
TSCKSUM
);
}
}
if
(
taosArrayPush
(
aBlockCol
,
pBlockCol
)
==
NULL
)
{
code
=
TSDB_CODE_OUT_OF_MEMORY
;
goto
_err
;
}
}
// write
code
=
tsdbWriteBlockDataImpl
(
pFileFD
,
pSubBlock
,
hdr
,
aBlockCol
,
*
ppBuf1
,
nData
,
ppBuf2
);
if
(
code
)
goto
_err
;
pSubBlock
->
szBlock
=
pSubBlock
->
szBlockCol
+
sizeof
(
TSCKSUM
)
+
nData
;
if
(
pBlock
->
last
)
{
if
(
pBlock
->
last
)
{
pWriter
->
wSet
.
fLast
.
size
+=
pSubBlock
->
szBlock
;
pWriter
->
wSet
.
fLast
.
size
+=
pSubBlock
->
szBlock
;
}
else
{
}
else
{
pWriter
->
wSet
.
fData
.
size
+=
pSubBlock
->
szBlock
;
pWriter
->
wSet
.
fData
.
size
+=
pSubBlock
->
szBlock
;
}
}
// ======================= BLOCK SMA =======================
pSubBlock
->
sOffset
=
0
;
pSubBlock
->
nSma
=
0
;
if
(
pBlock
->
nSubBlock
>
1
||
pBlock
->
last
||
pBlock
->
hasDup
)
goto
_exit
;
code
=
tsdbWriteBlockSMA
(
pWriter
->
pSmaFD
,
pBlockData
,
pSubBlock
,
ppBuf1
);
if
(
code
)
goto
_err
;
if
(
pSubBlock
->
nSma
>
0
)
{
pSubBlock
->
sOffset
=
pWriter
->
wSet
.
fSma
.
size
;
pWriter
->
wSet
.
fSma
.
size
+=
(
sizeof
(
SColumnDataAgg
)
*
pSubBlock
->
nSma
+
sizeof
(
TSCKSUM
));
}
_exit:
tsdbFree
(
pBuf1
);
tsdbFree
(
pBuf1
);
tsdbFree
(
pBuf2
);
tsdbFree
(
pBuf2
);
taosArrayDestroy
(
aBlockCol
);
return
code
;
return
code
;
_err:
_err:
tsdbError
(
"vgId:%d write block data failed since %s"
,
TD_VID
(
pWriter
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbError
(
"vgId:%d write block data failed since %s"
,
TD_VID
(
pWriter
->
pTsdb
->
pVnode
),
tstrerror
(
code
));
tsdbFree
(
pBuf1
);
tsdbFree
(
pBuf1
);
tsdbFree
(
pBuf2
);
tsdbFree
(
pBuf2
);
return
code
;
taosArrayDestroy
(
aBlockCol
);
}
int32_t
tsdbWriteBlockSMA
(
SDataFWriter
*
pWriter
,
SBlockSMA
*
pBlockSMA
,
int64_t
*
rOffset
,
int64_t
*
rSize
)
{
int32_t
code
=
0
;
// TODO
return
code
;
return
code
;
}
}
source/dnode/vnode/src/tsdb/tsdbUtil.c
浏览文件 @
7d50bfcb
...
@@ -15,57 +15,15 @@
...
@@ -15,57 +15,15 @@
#include "tsdb.h"
#include "tsdb.h"
#define TSDB_OFFSET_I32 ((uint8_t)0)
#define TSDB_OFFSET_I16 ((uint8_t)1)
#define TSDB_OFFSET_I8 ((uint8_t)2)
// SMapData =======================================================================
// SMapData =======================================================================
void
tMapDataReset
(
SMapData
*
pMapData
)
{
void
tMapDataReset
(
SMapData
*
pMapData
)
{
pMapData
->
flag
=
TSDB_OFFSET_I32
;
pMapData
->
nItem
=
0
;
pMapData
->
nItem
=
0
;
pMapData
->
nData
=
0
;
pMapData
->
nData
=
0
;
}
}
void
tMapDataClear
(
SMapData
*
pMapData
)
{
void
tMapDataClear
(
SMapData
*
pMapData
)
{
if
(
pMapData
->
pBuf
)
{
tsdbFree
((
uint8_t
*
)
pMapData
->
aOffset
);
tsdbFree
(
pMapData
->
pBuf
);
}
else
{
tsdbFree
(
pMapData
->
pOfst
);
tsdbFree
(
pMapData
->
pData
);
tsdbFree
(
pMapData
->
pData
);
}
}
int32_t
tMapDataCopy
(
SMapData
*
pMapDataSrc
,
SMapData
*
pMapDataDest
)
{
int32_t
code
=
0
;
int32_t
size
;
pMapDataDest
->
nItem
=
pMapDataSrc
->
nItem
;
pMapDataDest
->
flag
=
pMapDataSrc
->
flag
;
switch
(
pMapDataDest
->
flag
)
{
case
TSDB_OFFSET_I32
:
size
=
sizeof
(
int32_t
)
*
pMapDataDest
->
nItem
;
break
;
case
TSDB_OFFSET_I16
:
size
=
sizeof
(
int16_t
)
*
pMapDataDest
->
nItem
;
break
;
case
TSDB_OFFSET_I8
:
size
=
sizeof
(
int8_t
)
*
pMapDataDest
->
nItem
;
break
;
default:
ASSERT
(
0
);
}
code
=
tsdbRealloc
(
&
pMapDataDest
->
pOfst
,
size
);
if
(
code
)
goto
_exit
;
memcpy
(
pMapDataDest
->
pOfst
,
pMapDataSrc
->
pOfst
,
size
);
pMapDataDest
->
nData
=
pMapDataSrc
->
nData
;
code
=
tsdbRealloc
(
&
pMapDataDest
->
pData
,
pMapDataDest
->
nData
);
if
(
code
)
goto
_exit
;
memcpy
(
pMapDataDest
->
pData
,
pMapDataSrc
->
pData
,
pMapDataDest
->
nData
);
_exit:
return
code
;
}
}
int32_t
tMapDataPutItem
(
SMapData
*
pMapData
,
void
*
pItem
,
int32_t
(
*
tPutItemFn
)(
uint8_t
*
,
void
*
))
{
int32_t
tMapDataPutItem
(
SMapData
*
pMapData
,
void
*
pItem
,
int32_t
(
*
tPutItemFn
)(
uint8_t
*
,
void
*
))
{
...
@@ -77,35 +35,19 @@ int32_t tMapDataPutItem(SMapData *pMapData, void *pItem, int32_t (*tPutItemFn)(u
...
@@ -77,35 +35,19 @@ int32_t tMapDataPutItem(SMapData *pMapData, void *pItem, int32_t (*tPutItemFn)(u
pMapData
->
nData
+=
tPutItemFn
(
NULL
,
pItem
);
pMapData
->
nData
+=
tPutItemFn
(
NULL
,
pItem
);
// alloc
// alloc
code
=
tsdbRealloc
(
&
pMapData
->
pOfs
t
,
sizeof
(
int32_t
)
*
pMapData
->
nItem
);
code
=
tsdbRealloc
(
(
uint8_t
**
)
&
pMapData
->
aOffse
t
,
sizeof
(
int32_t
)
*
pMapData
->
nItem
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
code
=
tsdbRealloc
(
&
pMapData
->
pData
,
pMapData
->
nData
);
code
=
tsdbRealloc
(
&
pMapData
->
pData
,
pMapData
->
nData
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
// put
// put
ASSERT
(
pMapData
->
flag
==
TSDB_OFFSET_I32
);
pMapData
->
aOffset
[
nItem
]
=
offset
;
((
int32_t
*
)
pMapData
->
pOfst
)[
nItem
]
=
offset
;
tPutItemFn
(
pMapData
->
pData
+
offset
,
pItem
);
tPutItemFn
(
pMapData
->
pData
+
offset
,
pItem
);
_err:
_err:
return
code
;
return
code
;
}
}
static
FORCE_INLINE
int32_t
tMapDataGetOffset
(
SMapData
*
pMapData
,
int32_t
idx
)
{
switch
(
pMapData
->
flag
)
{
case
TSDB_OFFSET_I8
:
return
((
int8_t
*
)
pMapData
->
pOfst
)[
idx
];
break
;
case
TSDB_OFFSET_I16
:
return
((
int16_t
*
)
pMapData
->
pOfst
)[
idx
];
break
;
case
TSDB_OFFSET_I32
:
return
((
int32_t
*
)
pMapData
->
pOfst
)[
idx
];
break
;
default:
ASSERT
(
0
);
}
}
int32_t
tMapDataSearch
(
SMapData
*
pMapData
,
void
*
pSearchItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
),
int32_t
tMapDataSearch
(
SMapData
*
pMapData
,
void
*
pSearchItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
),
int32_t
(
*
tItemCmprFn
)(
const
void
*
,
const
void
*
),
void
*
pItem
)
{
int32_t
(
*
tItemCmprFn
)(
const
void
*
,
const
void
*
),
void
*
pItem
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
...
@@ -135,58 +77,25 @@ _exit:
...
@@ -135,58 +77,25 @@ _exit:
return
code
;
return
code
;
}
}
int32_t
tMapDataGetItemByIdx
(
SMapData
*
pMapData
,
int32_t
idx
,
void
*
pItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
))
{
void
tMapDataGetItemByIdx
(
SMapData
*
pMapData
,
int32_t
idx
,
void
*
pItem
,
int32_t
(
*
tGetItemFn
)(
uint8_t
*
,
void
*
))
{
int32_t
code
=
0
;
ASSERT
(
idx
>=
0
&&
idx
<
pMapData
->
nItem
);
tGetItemFn
(
pMapData
->
pData
+
pMapData
->
aOffset
[
idx
],
pItem
);
if
(
idx
<
0
||
idx
>=
pMapData
->
nItem
)
{
code
=
TSDB_CODE_NOT_FOUND
;
goto
_exit
;
}
tGetItemFn
(
pMapData
->
pData
+
tMapDataGetOffset
(
pMapData
,
idx
),
pItem
);
_exit:
return
code
;
}
}
int32_t
tPutMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
)
{
int32_t
tPutMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
)
{
int32_t
n
=
0
;
int32_t
n
=
0
;
ASSERT
(
pMapData
->
flag
==
TSDB_OFFSET_I32
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pMapData
->
nItem
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pMapData
->
nItem
);
if
(
pMapData
->
nItem
)
{
if
(
pMapData
->
nItem
)
{
int32_t
maxOffset
=
tMapDataGetOffset
(
pMapData
,
pMapData
->
nItem
-
1
);
if
(
maxOffset
<=
INT8_MAX
)
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I8
);
if
(
p
)
{
for
(
int32_t
iItem
=
0
;
iItem
<
pMapData
->
nItem
;
iItem
++
)
{
for
(
int32_t
iItem
=
0
;
iItem
<
pMapData
->
nItem
;
iItem
++
)
{
n
+=
tPutI8
(
p
+
n
,
(
int8_t
)
tMapDataGetOffset
(
pMapData
,
iItem
));
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pMapData
->
aOffset
[
iItem
]);
}
}
else
{
n
=
n
+
sizeof
(
int8_t
)
*
pMapData
->
nItem
;
}
}
else
if
(
maxOffset
<=
INT16_MAX
)
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I16
);
if
(
p
)
{
for
(
int32_t
iItem
=
0
;
iItem
<
pMapData
->
nItem
;
iItem
++
)
{
n
+=
tPutI16
(
p
+
n
,
(
int16_t
)
tMapDataGetOffset
(
pMapData
,
iItem
));
}
}
else
{
n
=
n
+
sizeof
(
int16_t
)
*
pMapData
->
nItem
;
}
}
}
else
{
n
+=
tPutU8
(
p
?
p
+
n
:
p
,
TSDB_OFFSET_I32
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pMapData
->
nData
);
if
(
p
)
{
if
(
p
)
{
for
(
int32_t
iItem
=
0
;
iItem
<
pMapData
->
nItem
;
iItem
++
)
{
memcpy
(
p
+
n
,
pMapData
->
pData
,
pMapData
->
nData
);
n
+=
tPutI32
(
p
+
n
,
tMapDataGetOffset
(
pMapData
,
iItem
));
}
}
else
{
n
=
n
+
sizeof
(
int32_t
)
*
pMapData
->
nItem
;
}
}
}
n
+=
pMapData
->
nData
;
n
+=
tPutBinary
(
p
?
p
+
n
:
p
,
pMapData
->
pData
,
pMapData
->
nData
);
}
}
return
n
;
return
n
;
...
@@ -194,26 +103,22 @@ int32_t tPutMapData(uint8_t *p, SMapData *pMapData) {
...
@@ -194,26 +103,22 @@ int32_t tPutMapData(uint8_t *p, SMapData *pMapData) {
int32_t
tGetMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
)
{
int32_t
tGetMapData
(
uint8_t
*
p
,
SMapData
*
pMapData
)
{
int32_t
n
=
0
;
int32_t
n
=
0
;
int32_t
offset
;
tMapDataReset
(
pMapData
);
n
+=
tGetI32v
(
p
+
n
,
&
pMapData
->
nItem
);
n
+=
tGetI32v
(
p
+
n
,
&
pMapData
->
nItem
);
if
(
pMapData
->
nItem
)
{
if
(
pMapData
->
nItem
)
{
n
+=
tGetU8
(
p
+
n
,
&
pMapData
->
flag
);
if
(
tsdbRealloc
((
uint8_t
**
)
&
pMapData
->
aOffset
,
sizeof
(
int32_t
)
*
pMapData
->
nItem
))
return
-
1
;
pMapData
->
pOfst
=
p
+
n
;
switch
(
pMapData
->
flag
)
{
for
(
int32_t
iItem
=
0
;
iItem
<
pMapData
->
nItem
;
iItem
++
)
{
case
TSDB_OFFSET_I8
:
n
+=
tGetI32v
(
p
+
n
,
&
pMapData
->
aOffset
[
iItem
]);
n
=
n
+
sizeof
(
int8_t
)
*
pMapData
->
nItem
;
break
;
case
TSDB_OFFSET_I16
:
n
=
n
+
sizeof
(
int16_t
)
*
pMapData
->
nItem
;
break
;
case
TSDB_OFFSET_I32
:
n
=
n
+
sizeof
(
int32_t
)
*
pMapData
->
nItem
;
break
;
default:
ASSERT
(
0
);
}
}
n
+=
tGetBinary
(
p
+
n
,
&
pMapData
->
pData
,
&
pMapData
->
nData
);
n
+=
tGetI32v
(
p
+
n
,
&
pMapData
->
nData
);
if
(
tsdbRealloc
(
&
pMapData
->
pData
,
pMapData
->
nData
))
return
-
1
;
memcpy
(
pMapData
->
pData
,
p
+
n
,
pMapData
->
nData
);
n
+=
pMapData
->
nData
;
}
}
return
n
;
return
n
;
...
@@ -377,55 +282,8 @@ int32_t tCmprBlockIdx(void const *lhs, void const *rhs) {
...
@@ -377,55 +282,8 @@ int32_t tCmprBlockIdx(void const *lhs, void const *rhs) {
// SBlock ======================================================
// SBlock ======================================================
void
tBlockReset
(
SBlock
*
pBlock
)
{
void
tBlockReset
(
SBlock
*
pBlock
)
{
pBlock
->
minKey
=
TSDBKEY_MAX
;
*
pBlock
=
pBlock
->
maxKey
=
TSDBKEY_MIN
;
(
SBlock
){.
minKey
=
TSDBKEY_MAX
,
.
maxKey
=
TSDBKEY_MIN
,
.
minVersion
=
VERSION_MAX
,
.
maxVersion
=
VERSION_MIN
};
pBlock
->
minVersion
=
VERSION_MAX
;
pBlock
->
maxVersion
=
VERSION_MIN
;
pBlock
->
nRow
=
0
;
pBlock
->
last
=
-
1
;
pBlock
->
hasDup
=
0
;
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
TSDB_MAX_SUBBLOCKS
;
iSubBlock
++
)
{
pBlock
->
aSubBlock
[
iSubBlock
].
nRow
=
0
;
pBlock
->
aSubBlock
[
iSubBlock
].
cmprAlg
=
-
1
;
pBlock
->
aSubBlock
[
iSubBlock
].
offset
=
-
1
;
pBlock
->
aSubBlock
[
iSubBlock
].
szVersion
=
-
1
;
pBlock
->
aSubBlock
[
iSubBlock
].
szTSKEY
=
-
1
;
pBlock
->
aSubBlock
[
iSubBlock
].
szBlock
=
-
1
;
tMapDataReset
(
&
pBlock
->
aSubBlock
->
mBlockCol
);
}
pBlock
->
nSubBlock
=
0
;
}
void
tBlockClear
(
SBlock
*
pBlock
)
{
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
TSDB_MAX_SUBBLOCKS
;
iSubBlock
++
)
{
tMapDataClear
(
&
pBlock
->
aSubBlock
->
mBlockCol
);
}
}
int32_t
tBlockCopy
(
SBlock
*
pBlockSrc
,
SBlock
*
pBlockDest
)
{
int32_t
code
=
0
;
pBlockDest
->
minKey
=
pBlockSrc
->
minKey
;
pBlockDest
->
maxKey
=
pBlockSrc
->
maxKey
;
pBlockDest
->
minVersion
=
pBlockSrc
->
minVersion
;
pBlockDest
->
maxVersion
=
pBlockSrc
->
maxVersion
;
pBlockDest
->
nRow
=
pBlockSrc
->
nRow
;
pBlockDest
->
last
=
pBlockSrc
->
last
;
pBlockDest
->
hasDup
=
pBlockSrc
->
hasDup
;
pBlockDest
->
nSubBlock
=
pBlockSrc
->
nSubBlock
;
for
(
int32_t
iSubBlock
=
0
;
iSubBlock
<
pBlockSrc
->
nSubBlock
;
iSubBlock
++
)
{
pBlockDest
->
aSubBlock
[
iSubBlock
].
nRow
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
nRow
;
pBlockDest
->
aSubBlock
[
iSubBlock
].
cmprAlg
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
cmprAlg
;
pBlockDest
->
aSubBlock
[
iSubBlock
].
offset
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
offset
;
pBlockDest
->
aSubBlock
[
iSubBlock
].
szVersion
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
szVersion
;
pBlockDest
->
aSubBlock
[
iSubBlock
].
szTSKEY
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
szTSKEY
;
pBlockDest
->
aSubBlock
[
iSubBlock
].
szBlock
=
pBlockSrc
->
aSubBlock
[
iSubBlock
].
szBlock
;
code
=
tMapDataCopy
(
&
pBlockSrc
->
aSubBlock
[
iSubBlock
].
mBlockCol
,
&
pBlockDest
->
aSubBlock
[
iSubBlock
].
mBlockCol
);
if
(
code
)
goto
_exit
;
}
_exit:
return
code
;
}
}
int32_t
tPutBlock
(
uint8_t
*
p
,
void
*
ph
)
{
int32_t
tPutBlock
(
uint8_t
*
p
,
void
*
ph
)
{
...
@@ -441,13 +299,15 @@ int32_t tPutBlock(uint8_t *p, void *ph) {
...
@@ -441,13 +299,15 @@ int32_t tPutBlock(uint8_t *p, void *ph) {
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
hasDup
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
hasDup
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
nSubBlock
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
nSubBlock
);
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
n
+=
tPutI
64
v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
nRow
);
n
+=
tPutI
32
v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
nRow
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
cmprAlg
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
cmprAlg
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
offset
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
offset
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szVersion
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szBlockCol
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szTSKEY
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szVersion
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szBlock
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szTSKEY
);
n
+=
tPutMapData
(
p
?
p
+
n
:
p
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
mBlockCol
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
szBlock
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
sOffset
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlock
->
aSubBlock
[
iSubBlock
].
nSma
);
}
}
return
n
;
return
n
;
...
@@ -466,20 +326,21 @@ int32_t tGetBlock(uint8_t *p, void *ph) {
...
@@ -466,20 +326,21 @@ int32_t tGetBlock(uint8_t *p, void *ph) {
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
hasDup
);
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
hasDup
);
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
nSubBlock
);
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
nSubBlock
);
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
for
(
int8_t
iSubBlock
=
0
;
iSubBlock
<
pBlock
->
nSubBlock
;
iSubBlock
++
)
{
n
+=
tGetI
64
v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
nRow
);
n
+=
tGetI
32
v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
nRow
);
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
cmprAlg
);
n
+=
tGetI8
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
cmprAlg
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
offset
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
offset
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szVersion
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szBlockCol
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szTSKEY
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szVersion
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szBlock
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szTSKEY
);
n
+=
tGetMapData
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
mBlockCol
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
szBlock
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
sOffset
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlock
->
aSubBlock
[
iSubBlock
].
nSma
);
}
}
return
n
;
return
n
;
}
}
int32_t
tBlockCmprFn
(
const
void
*
p1
,
const
void
*
p2
)
{
int32_t
tBlockCmprFn
(
const
void
*
p1
,
const
void
*
p2
)
{
int32_t
c
;
SBlock
*
pBlock1
=
(
SBlock
*
)
p1
;
SBlock
*
pBlock1
=
(
SBlock
*
)
p1
;
SBlock
*
pBlock2
=
(
SBlock
*
)
p2
;
SBlock
*
pBlock2
=
(
SBlock
*
)
p2
;
...
@@ -504,14 +365,11 @@ int32_t tPutBlockCol(uint8_t *p, void *ph) {
...
@@ -504,14 +365,11 @@ int32_t tPutBlockCol(uint8_t *p, void *ph) {
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlockCol
->
flag
);
n
+=
tPutI8
(
p
?
p
+
n
:
p
,
pBlockCol
->
flag
);
if
(
pBlockCol
->
flag
!=
HAS_NULL
)
{
if
(
pBlockCol
->
flag
!=
HAS_NULL
)
{
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlockCol
->
offset
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlockCol
->
offset
);
if
(
pBlockCol
->
flag
!=
HAS_VALUE
)
{
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlockCol
->
szBitmap
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlockCol
->
bsize
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlockCol
->
szOffset
);
}
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlockCol
->
szValue
);
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlockCol
->
csize
);
n
+=
tPutI32v
(
p
?
p
+
n
:
p
,
pBlockCol
->
szOrigin
);
if
(
IS_VAR_DATA_TYPE
(
pBlockCol
->
type
))
{
n
+=
tPutI64v
(
p
?
p
+
n
:
p
,
pBlockCol
->
osize
);
}
}
}
return
n
;
return
n
;
...
@@ -528,18 +386,11 @@ int32_t tGetBlockCol(uint8_t *p, void *ph) {
...
@@ -528,18 +386,11 @@ int32_t tGetBlockCol(uint8_t *p, void *ph) {
ASSERT
(
pBlockCol
->
flag
&&
(
pBlockCol
->
flag
!=
HAS_NONE
));
ASSERT
(
pBlockCol
->
flag
&&
(
pBlockCol
->
flag
!=
HAS_NONE
));
if
(
pBlockCol
->
flag
!=
HAS_NULL
)
{
if
(
pBlockCol
->
flag
!=
HAS_NULL
)
{
n
+=
tGetI64v
(
p
+
n
,
&
pBlockCol
->
offset
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlockCol
->
offset
);
if
(
pBlockCol
->
flag
!=
HAS_VALUE
)
{
n
+=
tGetI32v
(
p
+
n
,
&
pBlockCol
->
szBitmap
);
n
+=
tGetI64v
(
p
+
n
,
&
pBlockCol
->
bsize
);
n
+=
tGetI32v
(
p
+
n
,
&
pBlockCol
->
szOffset
);
}
else
{
n
+=
tGetI32v
(
p
+
n
,
&
pBlockCol
->
szValue
);
pBlockCol
->
bsize
=
0
;
n
+=
tGetI32v
(
p
+
n
,
&
pBlockCol
->
szOrigin
);
}
n
+=
tGetI64v
(
p
+
n
,
&
pBlockCol
->
csize
);
if
(
IS_VAR_DATA_TYPE
(
pBlockCol
->
type
))
{
n
+=
tGetI64v
(
p
+
n
,
&
pBlockCol
->
osize
);
}
else
{
pBlockCol
->
osize
=
-
1
;
}
}
}
return
n
;
return
n
;
...
@@ -942,12 +793,12 @@ int32_t tGetKEYINFO(uint8_t *p, KEYINFO *pKeyInfo) {
...
@@ -942,12 +793,12 @@ int32_t tGetKEYINFO(uint8_t *p, KEYINFO *pKeyInfo) {
}
}
// SColData ========================================
// SColData ========================================
void
tColDataReset
(
SColData
*
pColData
,
int16_t
cid
,
int8_t
type
)
{
void
tColDataReset
(
SColData
*
pColData
,
int16_t
cid
,
int8_t
type
,
int8_t
smaOn
)
{
pColData
->
cid
=
cid
;
pColData
->
cid
=
cid
;
pColData
->
type
=
type
;
pColData
->
type
=
type
;
pColData
->
smaOn
=
smaOn
;
pColData
->
nVal
=
0
;
pColData
->
nVal
=
0
;
pColData
->
flag
=
0
;
pColData
->
flag
=
0
;
pColData
->
offsetValid
=
0
;
pColData
->
nData
=
0
;
pColData
->
nData
=
0
;
}
}
...
@@ -977,26 +828,35 @@ int32_t tColDataAppendValue(SColData *pColData, SColVal *pColVal) {
...
@@ -977,26 +828,35 @@ int32_t tColDataAppendValue(SColData *pColData, SColVal *pColVal) {
if
(
pColVal
->
isNone
)
{
if
(
pColVal
->
isNone
)
{
pColData
->
flag
|=
HAS_NONE
;
pColData
->
flag
|=
HAS_NONE
;
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
0
);
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
0
);
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
pValue
=
NULL
;
}
else
if
(
pColVal
->
isNull
)
{
}
else
if
(
pColVal
->
isNull
)
{
pColData
->
flag
|=
HAS_NULL
;
pColData
->
flag
|=
HAS_NULL
;
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
1
);
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
1
);
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
pValue
=
NULL
;
}
else
{
}
else
{
pColData
->
flag
|=
HAS_VALUE
;
pColData
->
flag
|=
HAS_VALUE
;
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
2
);
SET_BIT2
(
pColData
->
pBitMap
,
pColData
->
nVal
,
2
);
pValue
=
&
pColVal
->
value
;
pValue
=
&
pColVal
->
value
;
}
}
if
(
pValue
)
{
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
{
code
=
tsdbRealloc
(
&
pColData
->
pData
,
pColData
->
nData
+
tPutValue
(
NULL
,
&
pColVal
->
value
,
pColVal
->
type
));
// offset
code
=
tsdbRealloc
((
uint8_t
**
)
&
pColData
->
aOffset
,
sizeof
(
int32_t
)
*
(
pColData
->
nVal
+
1
));
if
(
code
)
goto
_exit
;
if
(
code
)
goto
_exit
;
pColData
->
aOffset
[
pColData
->
nVal
]
=
pColData
->
nData
;
pColData
->
nData
+=
tPutValue
(
pColData
->
pData
+
pColData
->
nData
,
&
pColVal
->
value
,
pColVal
->
type
);
// value
if
((
!
pColVal
->
isNone
)
&&
(
!
pColVal
->
isNull
))
{
code
=
tsdbRealloc
(
&
pColData
->
pData
,
pColData
->
nData
+
pColVal
->
value
.
nData
);
if
(
code
)
goto
_exit
;
memcpy
(
pColData
->
pData
+
pColData
->
nData
,
pColVal
->
value
.
pData
,
pColVal
->
value
.
nData
);
pColData
->
nData
+=
pColVal
->
value
.
nData
;
}
}
else
{
code
=
tsdbRealloc
(
&
pColData
->
pData
,
pColData
->
nData
+
tPutValue
(
NULL
,
pValue
,
pColVal
->
type
));
if
(
code
)
goto
_exit
;
pColData
->
nData
+=
tPutValue
(
pColData
->
pData
+
pColData
->
nData
,
pValue
,
pColVal
->
type
);
}
}
pColData
->
nVal
++
;
pColData
->
nVal
++
;
pColData
->
offsetValid
=
0
;
_exit:
_exit:
return
code
;
return
code
;
...
@@ -1004,56 +864,32 @@ _exit:
...
@@ -1004,56 +864,32 @@ _exit:
int32_t
tColDataCopy
(
SColData
*
pColDataSrc
,
SColData
*
pColDataDest
)
{
int32_t
tColDataCopy
(
SColData
*
pColDataSrc
,
SColData
*
pColDataDest
)
{
int32_t
code
=
0
;
int32_t
code
=
0
;
int32_t
size
;
pColDataDest
->
cid
=
pColData
Dest
->
cid
;
pColDataDest
->
cid
=
pColData
Src
->
cid
;
pColDataDest
->
type
=
pColData
Dest
->
type
;
pColDataDest
->
type
=
pColData
Src
->
type
;
pColDataDest
->
offsetValid
=
0
;
pColDataDest
->
smaOn
=
pColDataSrc
->
smaOn
;
pColDataDest
->
nVal
=
pColDataSrc
->
nVal
;
pColDataDest
->
nVal
=
pColDataSrc
->
nVal
;
pColDataDest
->
flag
=
pColDataSrc
->
flag
;
pColDataDest
->
flag
=
pColDataSrc
->
flag
;
if
(
pColDataSrc
->
flag
!=
HAS_NONE
&&
pColDataSrc
->
flag
!=
HAS_NULL
&&
pColDataSrc
->
flag
!=
HAS_VALUE
)
{
code
=
tsdbRealloc
(
&
pColDataDest
->
pBitMap
,
BIT2_SIZE
(
pColDataDest
->
nVal
));
if
(
code
)
goto
_exit
;
memcpy
(
pColDataDest
->
pBitMap
,
pColDataSrc
->
pBitMap
,
BIT2_SIZE
(
pColDataSrc
->
nVal
));
size
=
BIT2_SIZE
(
pColDataSrc
->
nVal
);
}
code
=
tsdbRealloc
(
&
pColDataDest
->
pBitMap
,
size
);
pColDataDest
->
nData
=
pColDataSrc
->
nData
;
code
=
tsdbRealloc
(
&
pColDataDest
->
pData
,
pColDataSrc
->
nData
);
if
(
code
)
goto
_exit
;
if
(
code
)
goto
_exit
;
memcpy
(
pColDataDest
->
pData
,
pColDataSrc
->
pData
,
pColDataSrc
->
nData
);
memcpy
(
pColDataDest
->
pBitMap
,
pColDataSrc
->
pBitMap
,
size
);
_exit:
return
code
;
}
static
int32_t
tColDataUpdateOffset
(
SColData
*
pColData
)
{
int32_t
code
=
0
;
SValue
value
;
ASSERT
(
pColData
->
nVal
>
0
);
if
(
IS_VAR_DATA_TYPE
(
pColDataDest
->
type
))
{
ASSERT
(
pColData
->
flag
);
size
=
sizeof
(
int32_t
)
*
pColDataSrc
->
nVal
;
ASSERT
(
IS_VAR_DATA_TYPE
(
pColData
->
type
));
if
((
pColData
->
flag
&
HAS_VALUE
))
{
code
=
tsdbRealloc
((
uint8_t
**
)
&
pColDataDest
->
aOffset
,
size
);
code
=
tsdbRealloc
((
uint8_t
**
)
&
pColData
->
aOffset
,
sizeof
(
int32_t
)
*
pColData
->
nVal
);
if
(
code
)
goto
_exit
;
if
(
code
)
goto
_exit
;
int32_t
offset
=
0
;
memcpy
(
pColDataDest
->
aOffset
,
pColDataSrc
->
aOffset
,
size
);
for
(
int32_t
iVal
=
0
;
iVal
<
pColData
->
nVal
;
iVal
++
)
{
if
(
pColData
->
flag
!=
HAS_VALUE
)
{
uint8_t
v
=
GET_BIT2
(
pColData
->
pBitMap
,
iVal
);
if
(
v
==
0
||
v
==
1
)
{
pColData
->
aOffset
[
iVal
]
=
-
1
;
continue
;
}
}
pColData
->
aOffset
[
iVal
]
=
offset
;
offset
+=
tGetValue
(
pColData
->
pData
+
offset
,
&
value
,
pColData
->
type
);
}
}
ASSERT
(
offset
==
pColData
->
nData
);
code
=
tsdbRealloc
(
&
pColDataDest
->
pData
,
pColDataSrc
->
nData
);
pColData
->
offsetValid
=
1
;
if
(
code
)
goto
_exit
;
}
pColDataDest
->
nData
=
pColDataSrc
->
nData
;
memcpy
(
pColDataDest
->
pData
,
pColDataSrc
->
pData
,
pColDataDest
->
nData
);
_exit:
_exit:
return
code
;
return
code
;
...
@@ -1085,11 +921,13 @@ int32_t tColDataGetValue(SColData *pColData, int32_t iVal, SColVal *pColVal) {
...
@@ -1085,11 +921,13 @@ int32_t tColDataGetValue(SColData *pColData, int32_t iVal, SColVal *pColVal) {
// get value
// get value
SValue
value
;
SValue
value
;
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
{
if
(
IS_VAR_DATA_TYPE
(
pColData
->
type
))
{
if
(
!
pColData
->
offsetValid
)
{
if
(
iVal
+
1
<
pColData
->
nVal
)
{
code
=
tColDataUpdateOffset
(
pColData
);
value
.
nData
=
pColData
->
aOffset
[
iVal
+
1
]
-
pColData
->
aOffset
[
iVal
];
if
(
code
)
goto
_exit
;
}
else
{
value
.
nData
=
pColData
->
nData
-
pColData
->
aOffset
[
iVal
];
}
}
tGetValue
(
pColData
->
pData
+
pColData
->
aOffset
[
iVal
],
&
value
,
pColData
->
type
);
value
.
pData
=
pColData
->
pData
+
pColData
->
aOffset
[
iVal
];
}
else
{
}
else
{
tGetValue
(
pColData
->
pData
+
tDataTypes
[
pColData
->
type
].
bytes
*
iVal
,
&
value
,
pColData
->
type
);
tGetValue
(
pColData
->
pData
+
tDataTypes
[
pColData
->
type
].
bytes
*
iVal
,
&
value
,
pColData
->
type
);
}
}
...
@@ -1210,7 +1048,7 @@ int32_t tBlockDataAppendRow(SBlockData *pBlockData, TSDBROW *pRow, STSchema *pTS
...
@@ -1210,7 +1048,7 @@ int32_t tBlockDataAppendRow(SBlockData *pBlockData, TSDBROW *pRow, STSchema *pTS
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
// append a NONE
// append a NONE
tColDataReset
(
pColData
,
pColVal
->
cid
,
pColVal
->
type
);
tColDataReset
(
pColData
,
pColVal
->
cid
,
pColVal
->
type
,
0
);
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NONE
(
pColVal
->
cid
,
pColVal
->
type
));
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NONE
(
pColVal
->
cid
,
pColVal
->
type
));
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
@@ -1240,7 +1078,7 @@ int32_t tBlockDataAppendRow(SBlockData *pBlockData, TSDBROW *pRow, STSchema *pTS
...
@@ -1240,7 +1078,7 @@ int32_t tBlockDataAppendRow(SBlockData *pBlockData, TSDBROW *pRow, STSchema *pTS
code
=
tBlockDataAddColData
(
pBlockData
,
iColData
,
&
pColData
);
code
=
tBlockDataAddColData
(
pBlockData
,
iColData
,
&
pColData
);
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
tColDataReset
(
pColData
,
pColVal
->
cid
,
pColVal
->
type
);
tColDataReset
(
pColData
,
pColVal
->
cid
,
pColVal
->
type
,
0
);
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
for
(
int32_t
iRow
=
0
;
iRow
<
pBlockData
->
nRow
;
iRow
++
)
{
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NONE
(
pColVal
->
cid
,
pColVal
->
type
));
code
=
tColDataAppendValue
(
pColData
,
&
COL_VAL_NONE
(
pColVal
->
cid
,
pColVal
->
type
));
if
(
code
)
goto
_err
;
if
(
code
)
goto
_err
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录