Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
3c4b91a7
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1187
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看板
提交
3c4b91a7
编写于
5月 19, 2022
作者:
H
Hongze Cheng
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix: tdb concurrent w/r
上级
f3464aa0
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
129 addition
and
1 deletion
+129
-1
source/libs/tdb/src/db/tdbPager.c
source/libs/tdb/src/db/tdbPager.c
+2
-1
source/libs/tdb/test/tdbTest.cpp
source/libs/tdb/test/tdbTest.cpp
+127
-0
未找到文件。
source/libs/tdb/src/db/tdbPager.c
浏览文件 @
3c4b91a7
...
...
@@ -214,6 +214,8 @@ int tdbPagerCommit(SPager *pPager, TXN *pTxn) {
}
}
pPager
->
dbOrigSize
=
pPager
->
dbFileSize
;
// release the page
for
(
pPage
=
pPager
->
pDirty
;
pPage
;
pPage
=
pPager
->
pDirty
)
{
pPager
->
pDirty
=
pPage
->
pDirtyNext
;
...
...
@@ -230,7 +232,6 @@ int tdbPagerCommit(SPager *pPager, TXN *pTxn) {
// remote the journal file
tdbOsClose
(
pPager
->
jfd
);
tdbOsRemove
(
pPager
->
jFileName
);
pPager
->
dbOrigSize
=
pPager
->
dbFileSize
;
pPager
->
inTran
=
0
;
return
0
;
...
...
source/libs/tdb/test/tdbTest.cpp
浏览文件 @
3c4b91a7
#include <gtest/gtest.h>
#define ALLOW_FORBID_FUNC
#include "os.h"
#include "tdb.h"
#include <string>
#include <thread>
#include <vector>
typedef
struct
SPoolMem
{
int64_t
size
;
...
...
@@ -480,4 +483,128 @@ TEST(tdb_test, simple_upsert1) {
tdbTbClose
(
pDb
);
tdbClose
(
pEnv
);
}
TEST
(
tdb_test
,
multi_thread_query
)
{
int
ret
;
TDB
*
pEnv
;
TTB
*
pDb
;
tdb_cmpr_fn_t
compFunc
;
int
nData
=
20000
;
TXN
txn
;
taosRemoveDir
(
"tdb"
);
// Open Env
ret
=
tdbOpen
(
"tdb"
,
512
,
1
,
&
pEnv
);
GTEST_ASSERT_EQ
(
ret
,
0
);
// Create a database
compFunc
=
tKeyCmpr
;
ret
=
tdbTbOpen
(
"db.db"
,
-
1
,
-
1
,
compFunc
,
pEnv
,
&
pDb
);
GTEST_ASSERT_EQ
(
ret
,
0
);
char
key
[
64
];
char
val
[
64
];
int64_t
poolLimit
=
4096
;
// 1M pool limit
int64_t
txnid
=
0
;
SPoolMem
*
pPool
;
// open the pool
pPool
=
openPool
();
// start a transaction
txnid
++
;
txn
=
{.
flags
=
TDB_TXN_WRITE
|
TDB_TXN_READ_UNCOMMITTED
,
.
txnId
=
-
1
,
.
xMalloc
=
poolMalloc
,
.
xFree
=
poolFree
,
.
xArg
=
pPool
};
// tdbTxnOpen(&txn, txnid, poolMalloc, poolFree, pPool, );
tdbBegin
(
pEnv
,
&
txn
);
for
(
int
iData
=
1
;
iData
<=
nData
;
iData
++
)
{
sprintf
(
key
,
"key%d"
,
iData
);
sprintf
(
val
,
"value%d"
,
iData
);
ret
=
tdbTbInsert
(
pDb
,
key
,
strlen
(
key
),
val
,
strlen
(
val
),
&
txn
);
GTEST_ASSERT_EQ
(
ret
,
0
);
// if pool is full, commit the transaction and start a new one
// if (pPool->size >= poolLimit) {
// break;
// // commit current transaction
// tdbCommit(pEnv, &txn);
// tdbTxnClose(&txn);
// // start a new transaction
// clearPool(pPool);
// txnid++;
// tdbTxnOpen(&txn, txnid, poolMalloc, poolFree, pPool, TDB_TXN_WRITE | TDB_TXN_READ_UNCOMMITTED);
// tdbBegin(pEnv, &txn);
// }
}
auto
f
=
[](
TTB
*
pDb
,
int
nData
)
{
TBC
*
pDBC
;
void
*
pKey
=
NULL
;
void
*
pVal
=
NULL
;
int
vLen
,
kLen
;
int
count
=
0
;
int
ret
;
TXN
txn
;
SPoolMem
*
pPool
=
openPool
();
txn
=
{.
flags
=
0
,
.
txnId
=
0
,
.
xMalloc
=
poolMalloc
,
.
xFree
=
poolFree
,
.
xArg
=
pPool
};
ret
=
tdbTbcOpen
(
pDb
,
&
pDBC
,
&
txn
);
GTEST_ASSERT_EQ
(
ret
,
0
);
tdbTbcMoveToFirst
(
pDBC
);
for
(;;)
{
ret
=
tdbTbcNext
(
pDBC
,
&
pKey
,
&
kLen
,
&
pVal
,
&
vLen
);
if
(
ret
<
0
)
break
;
// std::cout.write((char *)pKey, kLen) /* << " " << kLen */ << " ";
// std::cout.write((char *)pVal, vLen) /* << " " << vLen */;
// std::cout << std::endl;
count
++
;
}
GTEST_ASSERT_EQ
(
count
,
nData
);
tdbTbcClose
(
pDBC
);
tdbFree
(
pKey
);
tdbFree
(
pVal
);
};
// tdbCommit(pEnv, &txn);
// multi-thread query
int
nThreads
=
20
;
std
::
vector
<
std
::
thread
>
threads
;
for
(
int
i
=
0
;
i
<
nThreads
;
i
++
)
{
if
(
i
==
0
)
{
threads
.
push_back
(
std
::
thread
(
tdbCommit
,
pEnv
,
&
txn
));
}
else
{
threads
.
push_back
(
std
::
thread
(
f
,
pDb
,
nData
));
}
}
for
(
auto
&
th
:
threads
)
{
th
.
join
();
}
// commit the transaction
tdbCommit
(
pEnv
,
&
txn
);
tdbTxnClose
(
&
txn
);
// Close a database
tdbTbClose
(
pDb
);
// Close Env
ret
=
tdbClose
(
pEnv
);
GTEST_ASSERT_EQ
(
ret
,
0
);
}
\ No newline at end of file
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录