Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
1ac9ff08
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看板
未验证
提交
1ac9ff08
编写于
9月 19, 2022
作者:
S
Shengliang Guan
提交者:
GitHub
9月 19, 2022
浏览文件
操作
浏览文件
下载
差异文件
Merge pull request #16926 from taosdata/feature/stream
enh: stream backend for sma
上级
fc0418fc
1604d729
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
13 addition
and
6 deletion
+13
-6
include/libs/stream/streamState.h
include/libs/stream/streamState.h
+1
-1
source/dnode/vnode/src/tq/tq.c
source/dnode/vnode/src/tq/tq.c
+2
-2
source/libs/stream/src/streamState.c
source/libs/stream/src/streamState.c
+7
-2
source/libs/wal/src/walMeta.c
source/libs/wal/src/walMeta.c
+3
-1
未找到文件。
include/libs/stream/streamState.h
浏览文件 @
1ac9ff08
...
...
@@ -34,7 +34,7 @@ typedef struct {
TXN
txn
;
}
SStreamState
;
SStreamState
*
streamStateOpen
(
char
*
path
,
SStreamTask
*
pTask
);
SStreamState
*
streamStateOpen
(
char
*
path
,
SStreamTask
*
pTask
,
bool
specPath
);
void
streamStateClose
(
SStreamState
*
pState
);
int32_t
streamStateBegin
(
SStreamState
*
pState
);
int32_t
streamStateCommit
(
SStreamState
*
pState
);
...
...
source/dnode/vnode/src/tq/tq.c
浏览文件 @
1ac9ff08
...
...
@@ -760,7 +760,7 @@ int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask) {
// expand executor
if
(
pTask
->
taskLevel
==
TASK_LEVEL__SOURCE
)
{
pTask
->
pState
=
streamStateOpen
(
pTq
->
pStreamMeta
->
path
,
pTask
);
pTask
->
pState
=
streamStateOpen
(
pTq
->
pStreamMeta
->
path
,
pTask
,
false
);
if
(
pTask
->
pState
==
NULL
)
{
return
-
1
;
}
...
...
@@ -774,7 +774,7 @@ int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask) {
pTask
->
exec
.
executor
=
qCreateStreamExecTaskInfo
(
pTask
->
exec
.
qmsg
,
&
handle
);
ASSERT
(
pTask
->
exec
.
executor
);
}
else
if
(
pTask
->
taskLevel
==
TASK_LEVEL__AGG
)
{
pTask
->
pState
=
streamStateOpen
(
pTq
->
pStreamMeta
->
path
,
pTask
);
pTask
->
pState
=
streamStateOpen
(
pTq
->
pStreamMeta
->
path
,
pTask
,
false
);
if
(
pTask
->
pState
==
NULL
)
{
return
-
1
;
}
...
...
source/libs/stream/src/streamState.c
浏览文件 @
1ac9ff08
...
...
@@ -18,14 +18,19 @@
#include "tcommon.h"
#include "ttimer.h"
SStreamState
*
streamStateOpen
(
char
*
path
,
SStreamTask
*
pTask
)
{
SStreamState
*
streamStateOpen
(
char
*
path
,
SStreamTask
*
pTask
,
bool
specPath
)
{
SStreamState
*
pState
=
taosMemoryCalloc
(
1
,
sizeof
(
SStreamState
));
if
(
pState
==
NULL
)
{
terrno
=
TSDB_CODE_OUT_OF_MEMORY
;
return
NULL
;
}
char
statePath
[
300
];
sprintf
(
statePath
,
"%s/%d"
,
path
,
pTask
->
taskId
);
if
(
!
specPath
)
{
sprintf
(
statePath
,
"%s/%d"
,
path
,
pTask
->
taskId
);
}
else
{
memcpy
(
statePath
,
path
,
300
);
}
if
(
tdbOpen
(
statePath
,
4096
,
256
,
&
pState
->
db
)
<
0
)
{
goto
_err
;
}
...
...
source/libs/wal/src/walMeta.c
浏览文件 @
1ac9ff08
...
...
@@ -268,7 +268,7 @@ int walRollFileInfo(SWal* pWal) {
char
*
walMetaSerialize
(
SWal
*
pWal
)
{
char
buf
[
30
];
ASSERT
(
pWal
->
fileInfoSet
);
int
sz
=
pWal
->
fileInfoSet
->
size
;
int
sz
=
taosArrayGetSize
(
pWal
->
fileInfoSet
)
;
cJSON
*
pRoot
=
cJSON_CreateObject
();
cJSON
*
pMeta
=
cJSON_CreateObject
();
cJSON
*
pFiles
=
cJSON_CreateArray
();
...
...
@@ -384,8 +384,10 @@ static int walFindCurMetaVer(SWal* pWal) {
int
code
=
regexec
(
&
walMetaRegexPattern
,
name
,
0
,
NULL
,
0
);
if
(
code
==
0
)
{
sscanf
(
name
,
"meta-ver%d"
,
&
metaVer
);
wDebug
(
"vgId:%d, wal find current meta: %s is the meta file, ver %d"
,
pWal
->
cfg
.
vgId
,
name
,
metaVer
);
break
;
}
wDebug
(
"vgId:%d, wal find current meta: %s is not meta file"
,
pWal
->
cfg
.
vgId
,
name
);
}
taosCloseDir
(
&
pDir
);
regfree
(
&
walMetaRegexPattern
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录