Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
慢慢CG
TDengine
提交
e35d89a4
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看板
提交
e35d89a4
编写于
7月 07, 2020
作者:
S
Shengliang Guan
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[TD-860] add processed count for sdb sync
上级
7fc487b8
变更
2
显示空白变更内容
内联
并排
Showing
2 changed file
with
28 addition
and
27 deletion
+28
-27
src/mnode/inc/mnodeSdb.h
src/mnode/inc/mnodeSdb.h
+1
-0
src/mnode/src/mnodeSdb.c
src/mnode/src/mnodeSdb.c
+27
-27
未找到文件。
src/mnode/inc/mnodeSdb.h
浏览文件 @
e35d89a4
...
...
@@ -53,6 +53,7 @@ typedef struct {
void
*
rowData
;
int32_t
rowSize
;
int32_t
retCode
;
// for callback in sdb queue
int32_t
processedCount
;
// for sync fwd callback
int32_t
(
*
cb
)(
struct
SMnodeMsg
*
pMsg
,
int32_t
code
);
struct
SMnodeMsg
*
pMsg
;
}
SSdbOper
;
...
...
src/mnode/src/mnodeSdb.c
浏览文件 @
e35d89a4
...
...
@@ -247,20 +247,22 @@ static void sdbConfirmForward(void *ahandle, void *param, int32_t code) {
assert
(
param
);
SSdbOper
*
pOper
=
param
;
SMnodeMsg
*
pMsg
=
pOper
->
pMsg
;
if
(
code
<=
0
)
pOper
->
retCode
=
code
;
if
(
code
>
0
)
{
int32_t
processedCount
=
atomic_add_fetch_32
(
&
pOper
->
processedCount
,
1
);
if
(
processedCount
<=
1
)
{
if
(
pMsg
!=
NULL
)
{
sdbDebug
(
"app:%p:%p, waiting for
slave to confirm this operation"
,
pMsg
->
rpcMsg
.
ahandle
,
pMsg
);
sdbDebug
(
"app:%p:%p, waiting for
confirm this operation, count:%d"
,
pMsg
->
rpcMsg
.
ahandle
,
pMsg
,
processedCount
);
}
return
;
}
if
(
pMsg
!=
NULL
)
{
sdbDebug
(
"app:%p:%p, is confirmed and will do callback func
, code:%s"
,
pMsg
->
rpcMsg
.
ahandle
,
pMsg
,
tstrerror
(
code
)
);
sdbDebug
(
"app:%p:%p, is confirmed and will do callback func
"
,
pMsg
->
rpcMsg
.
ahandle
,
pMsg
);
}
if
(
pOper
->
cb
!=
NULL
)
{
code
=
(
*
pOper
->
cb
)(
pMsg
,
c
ode
);
code
=
(
*
pOper
->
cb
)(
pMsg
,
pOper
->
retC
ode
);
}
dnodeSendRpcMnodeWriteRsp
(
pMsg
,
code
);
...
...
@@ -544,29 +546,32 @@ static int sdbWrite(void *param, void *data, int type) {
return
code
;
}
// forward to peers, even it is WAL/FWD, it shall be called to update version in sync
int32_t
syncCode
=
syncForwardToPeer
(
tsSdbObj
.
sync
,
pHead
,
pOper
,
TAOS_QTYPE_RPC
);
pthread_mutex_unlock
(
&
tsSdbObj
.
mutex
);
// from app, oper is created
if
(
pOper
!=
NULL
)
{
// forward to peers
int32_t
syncCode
=
syncForwardToPeer
(
tsSdbObj
.
sync
,
pHead
,
pOper
,
TAOS_QTYPE_RPC
);
if
(
syncCode
<=
0
)
atomic_add_fetch_32
(
&
pOper
->
processedCount
,
1
);
if
(
syncCode
<
0
)
{
sdbDebug
(
"table:%s, failed to forward request, result:%s action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
sdbError
(
"table:%s, failed to forward request, result:%s action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
tstrerror
(
syncCode
),
sdbGetActionStr
(
action
),
sdbGetKeyStr
(
pTable
,
pHead
->
cont
),
pHead
->
version
);
return
syncCode
;
}
else
if
(
syncCode
>
0
)
{
sdbDebug
(
"table:%s, forward request is sent, syncCode:%d
action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
syncCode
,
sdbGetActionStr
(
action
),
sdbGetKeyStr
(
pTable
,
pHead
->
cont
),
pHead
->
version
);
sdbDebug
(
"table:%s, forward request is sent,
action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
sdbGetActionStr
(
action
),
sdbGetKeyStr
(
pTable
,
pHead
->
cont
),
pHead
->
version
);
}
else
{
}
// from app, oper is created
if
(
pOper
!=
NULL
)
{
sdbDebug
(
"table:%s, record from app is disposed, action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
sdbTrace
(
"table:%s, no need to send fwd request, action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
sdbGetActionStr
(
action
),
sdbGetKeyStr
(
pTable
,
pHead
->
cont
),
pHead
->
version
);
}
return
syncCode
;
}
else
{
}
sdbDebug
(
"table:%s, record from wal/fwd is disposed, action:%s record:%s version:%"
PRId64
,
pTable
->
tableName
,
sdbGetActionStr
(
action
),
sdbGetKeyStr
(
pTable
,
pHead
->
cont
),
pHead
->
version
);
}
// even it is WAL/FWD, it shall be called to update version in sync
syncForwardToPeer
(
tsSdbObj
.
sync
,
pHead
,
pOper
,
TAOS_QTYPE_RPC
);
// from wal or forward msg, oper not created, should add into hash
if
(
action
==
SDB_ACTION_INSERT
)
{
...
...
@@ -972,11 +977,6 @@ static void *sdbWorkerFp(void *param) {
if
(
type
==
TAOS_QTYPE_RPC
)
{
pOper
=
(
SSdbOper
*
)
item
;
if
(
pOper
==
NULL
)
{
taosFreeQitem
(
item
);
continue
;
}
sdbDecRef
(
pOper
->
table
,
pOper
->
pObj
);
sdbConfirmForward
(
NULL
,
pOper
,
pOper
->
retCode
);
}
else
if
(
type
==
TAOS_QTYPE_FWD
)
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录