Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
287088ae
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看板
提交
287088ae
编写于
7月 26, 2023
作者:
dengyihao
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
fix stream state transfer
上级
c61393a7
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
26 addition
and
8 deletion
+26
-8
source/dnode/vnode/src/tq/tqStreamStateSnap.c
source/dnode/vnode/src/tq/tqStreamStateSnap.c
+5
-1
source/libs/stream/src/streamSnapshot.c
source/libs/stream/src/streamSnapshot.c
+21
-7
未找到文件。
source/dnode/vnode/src/tq/tqStreamStateSnap.c
浏览文件 @
287088ae
...
...
@@ -135,7 +135,9 @@ int32_t streamStateSnapWriterOpen(STQ* pTq, int64_t sver, int64_t ever, SStreamS
pWriter
->
sver
=
sver
;
pWriter
->
ever
=
ever
;
sprintf
(
tdir
,
"%s%s%s"
,
pTq
->
path
,
TD_DIRSEP
,
VNODE_TQ_STREAM
);
sprintf
(
tdir
,
"%s%s%s%s%s"
,
pTq
->
path
,
TD_DIRSEP
,
VNODE_TQ_STREAM
,
TD_DIRSEP
,
"received"
);
taosMkDir
(
tdir
);
SStreamSnapWriter
*
pSnapWriter
=
NULL
;
if
(
streamSnapWriterOpen
(
pTq
,
sver
,
ever
,
tdir
,
&
pSnapWriter
)
<
0
)
{
goto
_err
;
...
...
@@ -143,6 +145,8 @@ int32_t streamStateSnapWriterOpen(STQ* pTq, int64_t sver, int64_t ever, SStreamS
tqDebug
(
"vgId:%d, vnode stream-state snapshot writer opened, path:%s"
,
TD_VID
(
pTq
->
pVnode
),
tdir
);
pWriter
->
pWriterImpl
=
pSnapWriter
;
*
ppWriter
=
pWriter
;
return
code
;
_err:
tqError
(
"vgId:%d, vnode stream-state snapshot writer failed to open since %s"
,
TD_VID
(
pTq
->
pVnode
),
tstrerror
(
code
));
...
...
source/libs/stream/src/streamSnapshot.c
浏览文件 @
287088ae
...
...
@@ -361,8 +361,6 @@ int32_t streamSnapWriterOpen(void* pMeta, int64_t sver, int64_t ever, char* path
pHandle
->
currFileIdx
=
0
;
pHandle
->
offset
=
0
;
SBackendFileItem
*
pItem
=
taosArrayGet
(
pHandle
->
pFileList
,
pHandle
->
currFileIdx
);
pHandle
->
fd
=
streamOpenFile
(
pFile
->
path
,
pItem
->
name
,
TD_FILE_WRITE
);
*
ppWriter
=
pWriter
;
return
0
;
}
...
...
@@ -373,14 +371,25 @@ int32_t streamSnapWrite(SStreamSnapWriter* pWriter, uint8_t* pData, uint32_t nDa
SStreamSnapBlockHdr
*
pHdr
=
(
SStreamSnapBlockHdr
*
)
pData
;
SStreamSnapHandle
*
pHandle
=
&
pWriter
->
handle
;
SBanckendFile
*
pFile
=
pHandle
->
pBackendFile
;
SBackendFileItem
*
pItem
=
taosArrayGetP
(
pHandle
->
pFileList
,
pHandle
->
currFileIdx
);
SBackendFileItem
*
pItem
=
taosArrayGet
(
pHandle
->
pFileList
,
pHandle
->
currFileIdx
);
if
(
pHandle
->
fd
==
NULL
)
{
pHandle
->
fd
=
streamOpenFile
(
pFile
->
path
,
pItem
->
name
,
TD_FILE_CREATE
|
TD_FILE_WRITE
|
TD_FILE_APPEND
);
if
(
pHandle
->
fd
==
NULL
)
{
code
=
TAOS_SYSTEM_ERROR
(
terrno
);
qError
(
"stream-state failed to open file name:%s%s%s, reason:%s"
,
pFile
->
path
,
TD_DIRSEP
,
pHdr
->
name
,
tstrerror
(
code
));
}
}
if
(
strlen
(
pHdr
->
name
)
==
strlen
(
pItem
->
name
)
&&
strcmp
(
pHdr
->
name
,
pItem
->
name
)
==
0
)
{
if
(
taosPWriteFile
(
pHandle
->
fd
,
pHdr
->
data
,
pHdr
->
size
,
pHandle
->
offset
)
!=
pHdr
->
size
)
{
int64_t
bytes
=
taosPWriteFile
(
pHandle
->
fd
,
pHdr
->
data
,
pHdr
->
size
,
pHandle
->
offset
);
if
(
bytes
!=
pHdr
->
size
)
{
code
=
TAOS_SYSTEM_ERROR
(
terrno
);
qError
(
"stream
snap
failed to write snap, file name:%s, reason:%s"
,
pHdr
->
name
,
tstrerror
(
code
));
qError
(
"stream
-state
failed to write snap, file name:%s, reason:%s"
,
pHdr
->
name
,
tstrerror
(
code
));
return
code
;
}
pHandle
->
offset
+=
pHdr
->
size
;
pHandle
->
offset
+=
bytes
;
}
else
{
taosCloseFile
(
&
pHandle
->
fd
);
pHandle
->
offset
=
0
;
...
...
@@ -392,7 +401,12 @@ int32_t streamSnapWrite(SStreamSnapWriter* pWriter, uint8_t* pData, uint32_t nDa
taosArrayPush
(
pHandle
->
pFileList
,
&
item
);
SBackendFileItem
*
pItem
=
taosArrayGet
(
pHandle
->
pFileList
,
pHandle
->
currFileIdx
);
pHandle
->
fd
=
streamOpenFile
(
pFile
->
path
,
pItem
->
name
,
TD_FILE_WRITE
);
pHandle
->
fd
=
streamOpenFile
(
pFile
->
path
,
pItem
->
name
,
TD_FILE_CREATE
|
TD_FILE_WRITE
|
TD_FILE_APPEND
);
if
(
pHandle
->
fd
==
NULL
)
{
code
=
TAOS_SYSTEM_ERROR
(
terrno
);
qError
(
"stream-state failed to open file name:%s%s%s, reason:%s"
,
pFile
->
path
,
TD_DIRSEP
,
pHdr
->
name
,
tstrerror
(
code
));
}
taosPWriteFile
(
pHandle
->
fd
,
pHdr
->
data
,
pHdr
->
size
,
pHandle
->
offset
);
pHandle
->
offset
+=
pHdr
->
size
;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录