Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
weixin_44739061
DolphinScheduler
提交
86bcc058
DolphinScheduler
项目概览
weixin_44739061
/
DolphinScheduler
与 Fork 源项目一致
Fork自
apache / DolphinScheduler
通知
1
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,发现更多精彩内容 >>
未验证
提交
86bcc058
编写于
9月 04, 2020
作者:
C
CalvinKirs
提交者:
GitHub
9月 04, 2020
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[Test][remote]update netty util test and fix init channel error (#3671)
上级
d280820f
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
26 addition
and
19 deletion
+26
-19
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/NettyRemotingServer.java
...g/apache/dolphinscheduler/remote/NettyRemotingServer.java
+16
-17
dolphinscheduler-remote/src/test/java/org/apache/dolphinscheduler/remote/NettyUtilTest.java
...ava/org/apache/dolphinscheduler/remote/NettyUtilTest.java
+10
-2
未找到文件。
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/NettyRemotingServer.java
浏览文件 @
86bcc058
...
@@ -17,16 +17,6 @@
...
@@ -17,16 +17,6 @@
package
org.apache.dolphinscheduler.remote
;
package
org.apache.dolphinscheduler.remote
;
import
io.netty.bootstrap.ServerBootstrap
;
import
io.netty.channel.ChannelFuture
;
import
io.netty.channel.ChannelInitializer
;
import
io.netty.channel.ChannelOption
;
import
io.netty.channel.ChannelPipeline
;
import
io.netty.channel.EventLoopGroup
;
import
io.netty.channel.epoll.EpollEventLoopGroup
;
import
io.netty.channel.nio.NioEventLoopGroup
;
import
io.netty.channel.socket.nio.NioSocketChannel
;
import
org.apache.dolphinscheduler.remote.codec.NettyDecoder
;
import
org.apache.dolphinscheduler.remote.codec.NettyDecoder
;
import
org.apache.dolphinscheduler.remote.codec.NettyEncoder
;
import
org.apache.dolphinscheduler.remote.codec.NettyEncoder
;
import
org.apache.dolphinscheduler.remote.command.CommandType
;
import
org.apache.dolphinscheduler.remote.command.CommandType
;
...
@@ -36,15 +26,25 @@ import org.apache.dolphinscheduler.remote.processor.NettyRequestProcessor;
...
@@ -36,15 +26,25 @@ import org.apache.dolphinscheduler.remote.processor.NettyRequestProcessor;
import
org.apache.dolphinscheduler.remote.utils.Constants
;
import
org.apache.dolphinscheduler.remote.utils.Constants
;
import
org.apache.dolphinscheduler.remote.utils.NettyUtils
;
import
org.apache.dolphinscheduler.remote.utils.NettyUtils
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.ExecutorService
;
import
java.util.concurrent.Executors
;
import
java.util.concurrent.Executors
;
import
java.util.concurrent.ThreadFactory
;
import
java.util.concurrent.ThreadFactory
;
import
java.util.concurrent.atomic.AtomicBoolean
;
import
java.util.concurrent.atomic.AtomicBoolean
;
import
java.util.concurrent.atomic.AtomicInteger
;
import
java.util.concurrent.atomic.AtomicInteger
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
io.netty.bootstrap.ServerBootstrap
;
import
io.netty.channel.ChannelFuture
;
import
io.netty.channel.ChannelInitializer
;
import
io.netty.channel.ChannelOption
;
import
io.netty.channel.ChannelPipeline
;
import
io.netty.channel.EventLoopGroup
;
import
io.netty.channel.epoll.EpollEventLoopGroup
;
import
io.netty.channel.nio.NioEventLoopGroup
;
import
io.netty.channel.socket.SocketChannel
;
/**
/**
* remoting netty server
* remoting netty server
*/
*/
...
@@ -152,10 +152,10 @@ public class NettyRemotingServer {
...
@@ -152,10 +152,10 @@ public class NettyRemotingServer {
.
childOption
(
ChannelOption
.
TCP_NODELAY
,
serverConfig
.
isTcpNoDelay
())
.
childOption
(
ChannelOption
.
TCP_NODELAY
,
serverConfig
.
isTcpNoDelay
())
.
childOption
(
ChannelOption
.
SO_SNDBUF
,
serverConfig
.
getSendBufferSize
())
.
childOption
(
ChannelOption
.
SO_SNDBUF
,
serverConfig
.
getSendBufferSize
())
.
childOption
(
ChannelOption
.
SO_RCVBUF
,
serverConfig
.
getReceiveBufferSize
())
.
childOption
(
ChannelOption
.
SO_RCVBUF
,
serverConfig
.
getReceiveBufferSize
())
.
childHandler
(
new
ChannelInitializer
<
Nio
SocketChannel
>()
{
.
childHandler
(
new
ChannelInitializer
<
SocketChannel
>()
{
@Override
@Override
protected
void
initChannel
(
Nio
SocketChannel
ch
)
throws
Exception
{
protected
void
initChannel
(
SocketChannel
ch
)
throws
Exception
{
initNettyChannel
(
ch
);
initNettyChannel
(
ch
);
}
}
});
});
...
@@ -181,9 +181,8 @@ public class NettyRemotingServer {
...
@@ -181,9 +181,8 @@ public class NettyRemotingServer {
* init netty channel
* init netty channel
*
*
* @param ch socket channel
* @param ch socket channel
* @throws Exception
*/
*/
private
void
initNettyChannel
(
NioSocketChannel
ch
)
throws
Exception
{
private
void
initNettyChannel
(
SocketChannel
ch
)
{
ChannelPipeline
pipeline
=
ch
.
pipeline
();
ChannelPipeline
pipeline
=
ch
.
pipeline
();
pipeline
.
addLast
(
"encoder"
,
encoder
);
pipeline
.
addLast
(
"encoder"
,
encoder
);
pipeline
.
addLast
(
"decoder"
,
new
NettyDecoder
());
pipeline
.
addLast
(
"decoder"
,
new
NettyDecoder
());
...
...
dolphinscheduler-remote/src/test/java/org/apache/dolphinscheduler/remote/NettyUtilTest.java
浏览文件 @
86bcc058
...
@@ -17,20 +17,28 @@
...
@@ -17,20 +17,28 @@
package
org.apache.dolphinscheduler.remote
;
package
org.apache.dolphinscheduler.remote
;
import
static
org
.
apache
.
dolphinscheduler
.
remote
.
utils
.
Constants
.
OS_NAME
;
import
org.apache.dolphinscheduler.remote.utils.NettyUtils
;
import
org.apache.dolphinscheduler.remote.utils.NettyUtils
;
import
org.junit.Assert
;
import
org.junit.Assert
;
import
org.junit.Test
;
import
org.junit.Test
;
import
io.netty.channel.epoll.Epoll
;
/**
/**
* NettyUtilTest
* NettyUtilTest
*/
*/
public
class
NettyUtilTest
{
public
class
NettyUtilTest
{
@Test
@Test
public
void
testUserEpoll
()
{
public
void
testUserEpoll
()
{
System
.
setProperty
(
"netty.epoll.enable"
,
"false"
);
if
(
OS_NAME
.
toLowerCase
().
contains
(
"linux"
)
&&
Epoll
.
isAvailable
())
{
Assert
.
assertFalse
(
NettyUtils
.
useEpoll
());
Assert
.
assertTrue
(
NettyUtils
.
useEpoll
());
}
else
{
Assert
.
assertFalse
(
NettyUtils
.
useEpoll
());
}
}
}
}
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录