Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
4a259637
T
TDengine
项目概览
taosdata
/
TDengine
1 年多 前同步成功
通知
1185
Star
22016
Fork
4786
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
4a259637
编写于
11月 26, 2022
作者:
M
Minghao Li
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
refactor(sync): optimized heartbeat timer
上级
93efefcb
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
47 addition
and
26 deletion
+47
-26
source/libs/sync/inc/syncEnv.h
source/libs/sync/inc/syncEnv.h
+1
-0
source/libs/sync/inc/syncInt.h
source/libs/sync/inc/syncInt.h
+1
-0
source/libs/sync/src/syncMain.c
source/libs/sync/src/syncMain.c
+45
-26
未找到文件。
source/libs/sync/inc/syncEnv.h
浏览文件 @
4a259637
...
...
@@ -29,6 +29,7 @@ extern "C" {
#define ELECT_TIMER_MS_MAX (ELECT_TIMER_MS_MIN * 2)
#define ELECT_TIMER_MS_RANGE (ELECT_TIMER_MS_MAX - ELECT_TIMER_MS_MIN)
#define HEARTBEAT_TIMER_MS 1000
#define HEARTBEAT_TICK_NUM 20
typedef
struct
SSyncEnv
{
uint8_t
isStart
;
...
...
source/libs/sync/inc/syncInt.h
浏览文件 @
4a259637
...
...
@@ -61,6 +61,7 @@ typedef struct SSyncHbTimerData {
SSyncTimer
*
pTimer
;
SRaftId
destId
;
uint64_t
logicClock
;
int64_t
execTime
;
int64_t
rid
;
}
SSyncHbTimerData
;
...
...
source/libs/sync/src/syncMain.c
浏览文件 @
4a259637
...
...
@@ -711,9 +711,10 @@ static int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
pData
->
pTimer
=
pSyncTimer
;
pData
->
destId
=
pSyncTimer
->
destId
;
pData
->
logicClock
=
pSyncTimer
->
logicClock
;
pData
->
execTime
=
taosGetTimestampMs
()
+
pSyncTimer
->
timerMS
;
taosTmrReset
(
pSyncTimer
->
timerCb
,
pSyncTimer
->
timerMS
,
(
void
*
)(
pData
->
rid
),
syncEnv
()
->
pTimerManager
,
&
pSyncTimer
->
pTimer
);
taosTmrReset
(
pSyncTimer
->
timerCb
,
pSyncTimer
->
timerMS
/
HEARTBEAT_TICK_NUM
,
(
void
*
)(
pData
->
rid
)
,
syncEnv
()
->
pTimerManager
,
&
pSyncTimer
->
pTimer
);
}
else
{
sError
(
"vgId:%d, start ctrl hb timer error, sync env is stop"
,
pSyncNode
->
vgId
);
}
...
...
@@ -1979,6 +1980,7 @@ static void syncNodeEqHeartbeatTimer(void* param, void* tmrId) {
static
void
syncNodeEqPeerHeartbeatTimer
(
void
*
param
,
void
*
tmrId
)
{
int64_t
hbDataRid
=
(
int64_t
)
param
;
int64_t
tsNow
=
taosGetTimestampMs
();
SSyncHbTimerData
*
pData
=
syncHbTimerDataAcquire
(
hbDataRid
);
if
(
pData
==
NULL
)
{
...
...
@@ -2023,36 +2025,53 @@ static void syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId) {
int64_t
msgLogicClock
=
atomic_load_64
(
&
pData
->
logicClock
);
if
(
timerLogicClock
==
msgLogicClock
)
{
if
(
tsNow
>
pData
->
execTime
)
{
#if 0
sTrace(
"vgId:%d, hbDataRid:%ld, EXECUTE this step-------- heartbeat tsNow:%ld, exec:%ld, tsNow-exec:%ld, "
"---------",
pSyncNode->vgId, hbDataRid, tsNow, pData->execTime, tsNow - pData->execTime);
#endif
pData
->
execTime
+=
pSyncTimer
->
timerMS
;
SRpcMsg
rpcMsg
=
{
0
};
(
void
)
syncBuildHeartbeat
(
&
rpcMsg
,
pSyncNode
->
vgId
);
SyncHeartbeat
*
pSyncMsg
=
rpcMsg
.
pCont
;
pSyncMsg
->
srcId
=
pSyncNode
->
myRaftId
;
pSyncMsg
->
destId
=
pData
->
destId
;
pSyncMsg
->
term
=
pSyncNode
->
pRaftStore
->
currentTerm
;
pSyncMsg
->
commitIndex
=
pSyncNode
->
commitIndex
;
pSyncMsg
->
minMatchIndex
=
syncMinMatchIndex
(
pSyncNode
);
pSyncMsg
->
privateTerm
=
0
;
pSyncMsg
->
timeStamp
=
taosGetTimestampMs
();
// update reset time
int64_t
timerElapsed
=
tsNow
-
pSyncTimer
->
timeStamp
;
pSyncTimer
->
timeStamp
=
tsNow
;
char
logBuf
[
64
];
snprintf
(
logBuf
,
sizeof
(
logBuf
),
"timer-elapsed:%"
PRId64
", next-exec:%"
PRId64
,
timerElapsed
,
pData
->
execTime
);
// send msg
syncNodeSendHeartbeat
(
pSyncNode
,
&
pSyncMsg
->
destId
,
&
rpcMsg
,
logBuf
);
}
else
{
#if 0
sTrace(
"vgId:%d, hbDataRid:%ld, pass this step-------- heartbeat tsNow:%ld, exec:%ld, tsNow-exec:%ld, ---------",
pSyncNode->vgId, hbDataRid, tsNow, pData->execTime, tsNow - pData->execTime);
#endif
}
if
(
syncIsInit
())
{
// sTrace("vgId:%d, reset peer hb timer", pSyncNode->vgId);
taosTmrReset
(
syncNodeEqPeerHeartbeatTimer
,
pSyncTimer
->
timerMS
,
(
void
*
)
hbDataRid
,
syncEnv
()
->
pTimerManager
,
&
pSyncTimer
->
pTimer
);
taosTmrReset
(
syncNodeEqPeerHeartbeatTimer
,
pSyncTimer
->
timerMS
/
HEARTBEAT_TICK_NUM
,
(
void
*
)
hbDataRid
,
syncEnv
()
->
pTimerManager
,
&
pSyncTimer
->
pTimer
);
}
else
{
sError
(
"sync env is stop, reset peer hb timer error"
);
}
SRpcMsg
rpcMsg
=
{
0
};
(
void
)
syncBuildHeartbeat
(
&
rpcMsg
,
pSyncNode
->
vgId
);
SyncHeartbeat
*
pSyncMsg
=
rpcMsg
.
pCont
;
pSyncMsg
->
srcId
=
pSyncNode
->
myRaftId
;
pSyncMsg
->
destId
=
pData
->
destId
;
pSyncMsg
->
term
=
pSyncNode
->
pRaftStore
->
currentTerm
;
pSyncMsg
->
commitIndex
=
pSyncNode
->
commitIndex
;
pSyncMsg
->
minMatchIndex
=
syncMinMatchIndex
(
pSyncNode
);
pSyncMsg
->
privateTerm
=
0
;
pSyncMsg
->
timeStamp
=
taosGetTimestampMs
();
// update reset time
int64_t
tsNow
=
taosGetTimestampMs
();
int64_t
timerElapsed
=
tsNow
-
pSyncTimer
->
timeStamp
;
pSyncTimer
->
timeStamp
=
tsNow
;
char
logBuf
[
64
];
snprintf
(
logBuf
,
sizeof
(
logBuf
),
"timer-elapsed:%"
PRId64
,
timerElapsed
);
// send msg
syncNodeSendHeartbeat
(
pSyncNode
,
&
pSyncMsg
->
destId
,
&
rpcMsg
,
logBuf
);
}
else
{
sTrace
(
"vgId:%d, do not send hb, timerLogicClock:%"
PRId64
", msgLogicClock:%"
PRId64
""
,
pSyncNode
->
vgId
,
timerLogicClock
,
msgLogicClock
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录