Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
fdcdba62
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看板
未验证
提交
fdcdba62
编写于
11月 26, 2022
作者:
S
Shengliang Guan
提交者:
GitHub
11月 26, 2022
浏览文件
操作
浏览文件
下载
差异文件
Merge pull request #18472 from taosdata/feature/3.0_mhli
refactor(sync): optimized heartbeat timer
上级
38aa1a2f
b6dbc462
变更
6
隐藏空白更改
内联
并排
Showing
6 changed file
with
52 addition
and
31 deletion
+52
-31
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/inc/syncUtil.h
source/libs/sync/inc/syncUtil.h
+1
-1
source/libs/sync/src/syncMain.c
source/libs/sync/src/syncMain.c
+45
-26
source/libs/sync/src/syncReplication.c
source/libs/sync/src/syncReplication.c
+1
-1
source/libs/sync/src/syncUtil.c
source/libs/sync/src/syncUtil.c
+3
-3
未找到文件。
source/libs/sync/inc/syncEnv.h
浏览文件 @
fdcdba62
...
...
@@ -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
浏览文件 @
fdcdba62
...
...
@@ -61,6 +61,7 @@ typedef struct SSyncHbTimerData {
SSyncTimer
*
pTimer
;
SRaftId
destId
;
uint64_t
logicClock
;
int64_t
execTime
;
int64_t
rid
;
}
SSyncHbTimerData
;
...
...
source/libs/sync/inc/syncUtil.h
浏览文件 @
fdcdba62
...
...
@@ -94,7 +94,7 @@ void syncLogRecvLocalCmd(SSyncNode* pSyncNode, const SyncLocalCmd* pMsg, const c
void
syncLogSendAppendEntriesReply
(
SSyncNode
*
pSyncNode
,
const
SyncAppendEntriesReply
*
pMsg
,
const
char
*
s
);
void
syncLogRecvAppendEntriesReply
(
SSyncNode
*
pSyncNode
,
const
SyncAppendEntriesReply
*
pMsg
,
const
char
*
s
);
void
syncLogSendHeartbeat
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeat
*
pMsg
,
bool
printX
,
int64_t
timerElapsed
);
void
syncLogSendHeartbeat
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeat
*
pMsg
,
bool
printX
,
int64_t
timerElapsed
,
int64_t
execTime
);
void
syncLogRecvHeartbeat
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeat
*
pMsg
,
int64_t
timeDiff
);
void
syncLogSendHeartbeatReply
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeatReply
*
pMsg
,
const
char
*
s
);
...
...
source/libs/sync/src/syncMain.c
浏览文件 @
fdcdba62
...
...
@@ -698,6 +698,7 @@ static int32_t syncHbTimerInit(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer, SRa
static
int32_t
syncHbTimerStart
(
SSyncNode
*
pSyncNode
,
SSyncTimer
*
pSyncTimer
)
{
int32_t
ret
=
0
;
int64_t
tsNow
=
taosGetTimestampMs
();
if
(
syncIsInit
())
{
SSyncHbTimerData
*
pData
=
syncHbTimerDataAcquire
(
pSyncTimer
->
hbDataRid
);
if
(
pData
==
NULL
)
{
...
...
@@ -705,15 +706,16 @@ static int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
pData
->
rid
=
syncHbTimerDataAdd
(
pData
);
}
pSyncTimer
->
hbDataRid
=
pData
->
rid
;
pSyncTimer
->
timeStamp
=
t
aosGetTimestampMs
()
;
pSyncTimer
->
timeStamp
=
t
sNow
;
pData
->
syncNodeRid
=
pSyncNode
->
rid
;
pData
->
pTimer
=
pSyncTimer
;
pData
->
destId
=
pSyncTimer
->
destId
;
pData
->
logicClock
=
pSyncTimer
->
logicClock
;
pData
->
execTime
=
tsNow
+
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 +1981,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,35 +2026,51 @@ 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
=
tsNow
;
// update reset time
int64_t
timerElapsed
=
tsNow
-
pSyncTimer
->
timeStamp
;
pSyncTimer
->
timeStamp
=
tsNow
;
// send msg
syncLogSendHeartbeat
(
pSyncNode
,
pSyncMsg
,
false
,
timerElapsed
,
pData
->
execTime
);
syncNodeSendHeartbeat
(
pSyncNode
,
&
pSyncMsg
->
destId
,
&
rpcMsg
);
}
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
);
// update reset time
int64_t
tsNow
=
taosGetTimestampMs
();
int64_t
timerElapsed
=
tsNow
-
pSyncTimer
->
timeStamp
;
pSyncTimer
->
timeStamp
=
tsNow
;
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
=
tsNow
;
// send msg
syncLogSendHeartbeat
(
pSyncNode
,
pSyncMsg
,
false
,
timerElapsed
);
syncNodeSendHeartbeat
(
pSyncNode
,
&
pSyncMsg
->
destId
,
&
rpcMsg
);
}
else
{
sTrace
(
"vgId:%d, do not send hb, timerLogicClock:%"
PRId64
", msgLogicClock:%"
PRId64
""
,
pSyncNode
->
vgId
,
timerLogicClock
,
msgLogicClock
);
...
...
source/libs/sync/src/syncReplication.c
浏览文件 @
fdcdba62
...
...
@@ -230,7 +230,7 @@ int32_t syncNodeHeartbeatPeers(SSyncNode* pSyncNode) {
pSyncMsg
->
timeStamp
=
ts
;
// send msg
syncLogSendHeartbeat
(
pSyncNode
,
pSyncMsg
,
true
,
0
);
syncLogSendHeartbeat
(
pSyncNode
,
pSyncMsg
,
true
,
0
,
0
);
syncNodeSendHeartbeat
(
pSyncNode
,
&
pSyncMsg
->
destId
,
&
rpcMsg
);
}
...
...
source/libs/sync/src/syncUtil.c
浏览文件 @
fdcdba62
...
...
@@ -438,7 +438,7 @@ void syncLogRecvAppendEntriesReply(SSyncNode* pSyncNode, const SyncAppendEntries
host
,
port
,
pMsg
->
term
,
pMsg
->
privateTerm
,
pMsg
->
success
,
pMsg
->
lastSendIndex
,
pMsg
->
matchIndex
,
s
);
}
void
syncLogSendHeartbeat
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeat
*
pMsg
,
bool
printX
,
int64_t
timerElapsed
)
{
void
syncLogSendHeartbeat
(
SSyncNode
*
pSyncNode
,
const
SyncHeartbeat
*
pMsg
,
bool
printX
,
int64_t
timerElapsed
,
int64_t
execTime
)
{
if
(
!
(
sDebugFlag
&
DEBUG_TRACE
))
return
;
char
host
[
64
];
...
...
@@ -453,8 +453,8 @@ void syncLogSendHeartbeat(SSyncNode* pSyncNode, const SyncHeartbeat* pMsg, bool
}
else
{
sNTrace
(
pSyncNode
,
"send sync-heartbeat to %s:%d {term:%"
PRId64
", cmt:%"
PRId64
", min-match:%"
PRId64
", ts:%"
PRId64
"}, timer-elapsed:%"
PRId64
,
host
,
port
,
pMsg
->
term
,
pMsg
->
commitIndex
,
pMsg
->
minMatchIndex
,
pMsg
->
timeStamp
,
timerElapsed
);
"}, timer-elapsed:%"
PRId64
", next-exec:%"
PRId64
,
host
,
port
,
pMsg
->
term
,
pMsg
->
commitIndex
,
pMsg
->
minMatchIndex
,
pMsg
->
timeStamp
,
timerElapsed
,
execTime
);
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录