Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
liyuanzhong001
DolphinScheduler
提交
3191eb92
DolphinScheduler
项目概览
liyuanzhong001
/
DolphinScheduler
与 Fork 源项目一致
Fork自
apache / DolphinScheduler
通知
11
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
DolphinScheduler
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
提交
3191eb92
编写于
7月 25, 2019
作者:
leon-baoliang
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
update zkclient
上级
9361cbbf
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
30 addition
and
43 deletion
+30
-43
escheduler-common/src/main/java/cn/escheduler/common/zk/AbstractZKClient.java
...c/main/java/cn/escheduler/common/zk/AbstractZKClient.java
+27
-40
escheduler-server/src/main/java/cn/escheduler/server/zk/ZKMasterClient.java
...src/main/java/cn/escheduler/server/zk/ZKMasterClient.java
+3
-3
未找到文件。
escheduler-common/src/main/java/cn/escheduler/common/zk/AbstractZKClient.java
浏览文件 @
3191eb92
...
...
@@ -211,49 +211,11 @@ public abstract class AbstractZKClient {
return
false
;
}
/**
* init system znode
*/
protected
void
initSystemZNode
(){
try
{
// read master node parent path from conf
masterZNodeParentPath
=
getMasterZNodeParentPath
();
// read worker node parent path from conf
workerZNodeParentPath
=
getWorkerZNodeParentPath
();
// read server node parent path from conf
deadServerZNodeParentPath
=
getDeadZNodeParentPath
();
if
(
zkClient
.
checkExists
().
forPath
(
deadServerZNodeParentPath
)
==
null
){
// create persistent dead server parent node
zkClient
.
create
().
creatingParentContainersIfNeeded
()
.
withMode
(
CreateMode
.
PERSISTENT
).
forPath
(
deadServerZNodeParentPath
);
}
if
(
zkClient
.
checkExists
().
forPath
(
masterZNodeParentPath
)
==
null
){
// create persistent master parent node
zkClient
.
create
().
creatingParentContainersIfNeeded
()
.
withMode
(
CreateMode
.
PERSISTENT
).
forPath
(
masterZNodeParentPath
);
}
if
(
zkClient
.
checkExists
().
forPath
(
workerZNodeParentPath
)
==
null
){
// create persistent worker parent node
zkClient
.
create
().
creatingParentContainersIfNeeded
()
.
withMode
(
CreateMode
.
PERSISTENT
).
forPath
(
workerZNodeParentPath
);
}
}
catch
(
Exception
e
)
{
logger
.
error
(
"init system znode failed : "
+
e
.
getMessage
(),
e
);
}
}
public
void
removeDeadServerByHost
(
String
host
,
String
serverType
)
throws
Exception
{
List
<
String
>
deadServers
=
zkClient
.
getChildren
().
forPath
(
deadServerZNodeParentPath
);
for
(
String
serverPath
:
deadServers
){
if
(
serverPath
.
startsWith
(
serverType
+
UNDERLINE
+
host
)){
String
server
=
deadServerZNodeParentPath
+
SINGLE_SLASH
+
serverPath
;
zkClient
.
delete
().
forPath
(
server
);
logger
.
info
(
"{} server {} deleted from zk dead server path success"
,
serverType
,
host
);
...
...
@@ -394,10 +356,11 @@ public abstract class AbstractZKClient {
* @return
* @throws Exception
*/
public
boolean
checkZKNodeExists
(
String
host
,
ZKNodeType
zkNodeType
)
throws
Exception
{
public
boolean
checkZKNodeExists
(
String
host
,
ZKNodeType
zkNodeType
)
{
String
path
=
getZNodeParentPath
(
zkNodeType
);
if
(
StringUtils
.
isEmpty
(
path
)){
logger
.
error
(
"check zk node exists error, host:{}, zk node type:{}"
,
host
,
zkNodeType
.
toString
());
logger
.
error
(
"check zk node exists error, host:{}, zk node type:{}"
,
host
,
zkNodeType
.
toString
());
return
false
;
}
Map
<
String
,
String
>
serverMaps
=
getServerList
(
zkNodeType
);
...
...
@@ -508,7 +471,31 @@ public abstract class AbstractZKClient {
}
}
/**
* init system znode
*/
protected
void
initSystemZNode
(){
try
{
createNodePath
(
getMasterZNodeParentPath
());
createNodePath
(
getWorkerZNodeParentPath
());
createNodePath
(
getDeadZNodeParentPath
());
}
catch
(
Exception
e
)
{
logger
.
error
(
"init system znode failed : "
+
e
.
getMessage
(),
e
);
}
}
/**
* create zookeeper node path if not exists
* @param zNodeParentPath
* @throws Exception
*/
private
void
createNodePath
(
String
zNodeParentPath
)
throws
Exception
{
if
(
null
==
zkClient
.
checkExists
().
forPath
(
zNodeParentPath
)){
zkClient
.
create
().
creatingParentContainersIfNeeded
()
.
withMode
(
CreateMode
.
PERSISTENT
).
forPath
(
zNodeParentPath
);
}
}
@Override
public
String
toString
()
{
...
...
escheduler-server/src/main/java/cn/escheduler/server/zk/ZKMasterClient.java
浏览文件 @
3191eb92
...
...
@@ -322,9 +322,9 @@ public class ZKMasterClient extends AbstractZKClient {
// handle dead server
handleDeadServer
(
path
,
Constants
.
WORKER_PREFIX
,
Constants
.
ADD_ZK_OP
);
// create a distributed lock
, and the root node path of the lock space is /escheduler/lock/failover/worker
String
znodeLock
=
zkMasterClient
.
getWorkerFailoverLockPath
();
mutex
=
new
InterProcessMutex
(
zkMasterClient
.
getZkClient
(),
znodeLock
);
// create a distributed lock
String
znodeLock
=
getWorkerFailoverLockPath
();
mutex
=
new
InterProcessMutex
(
getZkClient
(),
znodeLock
);
mutex
.
acquire
();
String
workerHost
=
getHostByEventDataPath
(
path
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录