Browse Source

rpc

pull/3/MERGE
CalvinKirs 4 years ago
parent
commit
939ee1f2f5
  1. 7
      dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/MainTest.java
  2. 10
      dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/common/RpcRequest.java
  3. 2
      dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyClientHandler.java
  4. 6
      dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyServer.java
  5. 3
      dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyServerHandler.java

7
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/MainTest.java

@ -1,7 +1,6 @@
package org.apache.dolphinscheduler.remote.rpc;
import org.apache.dolphinscheduler.remote.config.NettyServerConfig;
import org.apache.dolphinscheduler.remote.rpc.client.IRpcClient;
import org.apache.dolphinscheduler.remote.rpc.client.RpcClient;
import org.apache.dolphinscheduler.remote.rpc.remote.NettyServer;
@ -14,6 +13,8 @@ import org.apache.dolphinscheduler.remote.utils.Host;
public class MainTest {
public static void main(String[] args) throws Exception {
NettyServer nettyServer = new NettyServer(new NettyServerConfig());
// NettyClient nettyClient=new NettyClient(new NettyClientConfig());
@ -31,9 +32,5 @@ public class MainTest {
UserService user = rpcClient.create(UserService.class, host);
System.out.println(user.hi(99999));
System.out.println(user.hi(998888888));
// UserCallback.class.newInstance().run("lllll");
// nettyClient.sendMsg(host,rpcRequest);
}
}

10
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/common/RpcRequest.java

@ -10,6 +10,16 @@ public class RpcRequest {
private String methodName;
private Class<?>[] parameterTypes;
private Object[] parameters;
// 0 hear beat,1 businness msg
private Byte eventType=1;
public Byte getEventType() {
return eventType;
}
public void setEventType(Byte eventType) {
this.eventType = eventType;
}
public String getRequestId() {
return requestId;

2
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyClientHandler.java

@ -76,7 +76,7 @@ public class NettyClientHandler extends ChannelInboundHandlerAdapter {
if (evt instanceof IdleStateEvent) {
IdleStateEvent event = (IdleStateEvent) evt;
RpcRequest request = new RpcRequest();
request.setMethodName("heart");
request.setEventType((byte)0);
ctx.channel().writeAndFlush(request);
logger.debug("send heart beat msg...");

6
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyServer.java

@ -128,7 +128,7 @@ public class NettyServer {
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
protected void initChannel(SocketChannel ch){
initNettyChannel(ch);
}
});
@ -137,11 +137,11 @@ public class NettyServer {
try {
future = serverBootstrap.bind(serverConfig.getListenPort()).sync();
} catch (Exception e) {
//logger.error("NettyRemotingServer bind fail {}, exit", e.getMessage(), e);
logger.error("NettyRemotingServer bind fail {}, exit", e.getMessage(), e);
throw new RuntimeException(String.format("NettyRemotingServer bind %s fail", serverConfig.getListenPort()));
}
if (future.isSuccess()) {
// logger.info("NettyRemotingServer bind success at port : {}", serverConfig.getListenPort());
logger.info("NettyRemotingServer bind success at port : {}", serverConfig.getListenPort());
} else if (future.cause() != null) {
throw new RuntimeException(String.format("NettyRemotingServer bind %s fail", serverConfig.getListenPort()), future.cause());
} else {

3
dolphinscheduler-remote/src/main/java/org/apache/dolphinscheduler/remote/rpc/remote/NettyServerHandler.java

@ -43,7 +43,8 @@ public class NettyServerHandler extends ChannelInboundHandlerAdapter {
RpcRequest req = (RpcRequest) msg;
RpcResponse response = new RpcResponse();
if (req.getMethodName().equals("heart")) {
if (req.getEventType() == 0) {
logger.info("接受心跳消息!...");
return;
}

Loading…
Cancel
Save