Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
B
byit-myth-job
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
liyuan
byit-myth-job
Commits
5feb60ca
Commit
5feb60ca
authored
Nov 05, 2020
by
huangfusuper
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
错误日志上传
parent
fa446771
Show whitespace changes
Inline
Side-by-side
Showing
7 changed files
with
99 additions
and
9 deletions
+99
-9
AdminRegisterConfig.java
...in/src/main/java/com/byit/config/AdminRegisterConfig.java
+6
-0
KeyUtil.java
...e/myth-dto-core/src/main/java/com/byit/utils/KeyUtil.java
+0
-0
MythJobProcess.java
...server/src/main/java/com/byit/process/MythJobProcess.java
+0
-1
pom.xml
byit-plugin-core/byit-plugin-rpc-common/pom.xml
+6
-0
PluginClientInitialization.java
...c/main/java/com/byit/init/PluginClientInitialization.java
+1
-1
ClientRpcUtil.java
...pc-client/src/main/java/com/byit/utils/ClientRpcUtil.java
+20
-7
RpcSpringUtil.java
...pc-client/src/main/java/com/byit/utils/RpcSpringUtil.java
+66
-0
No files found.
byit-myth-admin/src/main/java/com/byit/config/AdminRegisterConfig.java
View file @
5feb60ca
...
@@ -3,6 +3,7 @@ package com.byit.config;
...
@@ -3,6 +3,7 @@ package com.byit.config;
import
com.byit.rpc.registry.impl.RegistryServiceRegistry
;
import
com.byit.rpc.registry.impl.RegistryServiceRegistry
;
import
com.byit.rpc.remoting.invoker.impl.RpcSpringInvokerFactory
;
import
com.byit.rpc.remoting.invoker.impl.RpcSpringInvokerFactory
;
import
com.byit.rpc.remoting.provider.impl.RpcSpringProviderFactory
;
import
com.byit.rpc.remoting.provider.impl.RpcSpringProviderFactory
;
import
com.byit.utils.RpcSpringUtil
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.factory.annotation.Value
;
import
org.springframework.beans.factory.annotation.Value
;
import
org.springframework.context.annotation.Bean
;
import
org.springframework.context.annotation.Bean
;
...
@@ -62,4 +63,9 @@ public class AdminRegisterConfig {
...
@@ -62,4 +63,9 @@ public class AdminRegisterConfig {
return
providerFactory
;
return
providerFactory
;
}
}
@Bean
public
RpcSpringUtil
rpcSpringUtil
(){
return
new
RpcSpringUtil
();
}
}
}
byit-myth-
executor/myth-executor-server
/src/main/java/com/byit/utils/KeyUtil.java
→
byit-myth-
core/myth-dto-core
/src/main/java/com/byit/utils/KeyUtil.java
View file @
5feb60ca
File moved
byit-myth-executor/myth-executor-server/src/main/java/com/byit/process/MythJobProcess.java
View file @
5feb60ca
...
@@ -13,7 +13,6 @@ import org.springframework.data.redis.core.StringRedisTemplate;
...
@@ -13,7 +13,6 @@ import org.springframework.data.redis.core.StringRedisTemplate;
import
java.io.*
;
import
java.io.*
;
import
java.lang.reflect.Field
;
import
java.lang.reflect.Field
;
import
java.nio.channels.FileChannel
;
import
java.util.List
;
import
java.util.List
;
import
java.util.Map
;
import
java.util.Map
;
import
java.util.Objects
;
import
java.util.Objects
;
...
...
byit-plugin-core/byit-plugin-rpc-common/pom.xml
View file @
5feb60ca
...
@@ -15,6 +15,12 @@
...
@@ -15,6 +15,12 @@
<dependencies>
<dependencies>
<dependency>
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter-redis
</artifactId>
<version>
1.3.2.RELEASE
</version>
</dependency>
<dependency>
<groupId>
com.alibaba
</groupId>
<groupId>
com.alibaba
</groupId>
<artifactId>
fastjson
</artifactId>
<artifactId>
fastjson
</artifactId>
</dependency>
</dependency>
...
...
byit-plugin-core/myth-plugin-rpc-client/src/main/java/com/byit/init/PluginClientInitialization.java
View file @
5feb60ca
...
@@ -71,7 +71,7 @@ public class PluginClientInitialization {
...
@@ -71,7 +71,7 @@ public class PluginClientInitialization {
String
address
=
null
;
String
address
=
null
;
TreeSet
<
String
>
discovery
=
pluginServiceRegistry
.
discovery
(
pluginRpcRequestPacket
.
getJobName
());
TreeSet
<
String
>
discovery
=
pluginServiceRegistry
.
discovery
(
pluginRpcRequestPacket
.
getJobName
());
if
(
CollectionUtil
.
isNotEmpty
(
discovery
)){
if
(
CollectionUtil
.
isNotEmpty
(
discovery
)){
address
=
ClientRpcUtil
.
selectRPCServer
(
pluginClient
,
loadBalance
,
discovery
,
pluginRpcRequestPacket
.
getJobName
()
);
address
=
ClientRpcUtil
.
selectRPCServer
(
pluginClient
,
loadBalance
,
discovery
,
pluginRpcRequestPacket
);
// if(discovery.size() ==1){
// if(discovery.size() ==1){
// address = discovery.first();
// address = discovery.first();
// }else{
// }else{
...
...
byit-plugin-core/myth-plugin-rpc-client/src/main/java/com/byit/utils/ClientRpcUtil.java
View file @
5feb60ca
package
com
.
byit
.
utils
;
package
com
.
byit
.
utils
;
import
com.alibaba.fastjson.JSON
;
import
com.byit.client.PluginClient
;
import
com.byit.client.PluginClient
;
import
com.byit.dto.plugin.RunLog
;
import
com.byit.packet.request.PluginRpcRequestPacket
;
import
com.byit.param.PluginBeat
;
import
com.byit.param.PluginBeat
;
import
com.byit.rpc.remoting.invoker.route.LoadBalance
;
import
com.byit.rpc.remoting.invoker.route.LoadBalance
;
import
com.byit.rpc.util.RPCLogUtil
;
import
com.byit.rpc.util.RPCLogUtil
;
import
com.byit.rpc.util.RpcException
;
import
com.byit.rpc.util.RpcException
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.data.redis.core.StringRedisTemplate
;
import
java.util.TreeSet
;
import
java.util.TreeSet
;
import
static
com
.
alibaba
.
fastjson
.
serializer
.
SerializerFeature
.
WriteClassName
;
@Slf4j
@Slf4j
public
class
ClientRpcUtil
{
public
class
ClientRpcUtil
{
/**
/**
...
@@ -17,14 +23,16 @@ public class ClientRpcUtil {
...
@@ -17,14 +23,16 @@ public class ClientRpcUtil {
* @param client 客户端链接帮助其
* @param client 客户端链接帮助其
* @param loadBalance 负载均衡器
* @param loadBalance 负载均衡器
* @param address 地址
* @param address 地址
* @param
serverKey
服务名称
* @param
pluginRpcRequestPacket
服务名称
* @return 健全的服务地址
* @return 健全的服务地址
* @throws InterruptedException 线程异常
* @throws InterruptedException 线程异常
*/
*/
public
static
String
selectRPCServer
(
PluginClient
client
,
LoadBalance
loadBalance
,
TreeSet
<
String
>
address
,
String
serverKey
)
throws
InterruptedException
{
public
static
String
selectRPCServer
(
PluginClient
client
,
LoadBalance
loadBalance
,
TreeSet
<
String
>
address
,
PluginRpcRequestPacket
pluginRpcRequestPacket
)
throws
InterruptedException
{
for
(
String
ignored
:
address
)
{
for
(
String
ignored
:
address
)
{
String
routeHost
=
loadBalance
.
rpcInvokerRouter
.
route
(
serverKey
,
address
);
String
routeHost
=
loadBalance
.
rpcInvokerRouter
.
route
(
pluginRpcRequestPacket
.
getJobName
(),
address
);
boolean
retryRpcHost
=
retryRpcHost
(
client
,
routeHost
,
3
,
1
);
String
extension
=
pluginRpcRequestPacket
.
getExtension
();
String
runKey
=
KeyUtil
.
generateRunKey
(
Integer
.
parseInt
(
extension
));
boolean
retryRpcHost
=
retryRpcHost
(
client
,
routeHost
,
runKey
,
3
,
1
);
if
(
retryRpcHost
)
{
if
(
retryRpcHost
)
{
return
routeHost
;
return
routeHost
;
}
}
...
@@ -41,7 +49,7 @@ public class ClientRpcUtil {
...
@@ -41,7 +49,7 @@ public class ClientRpcUtil {
* @param thisRetryCount 当前重试次数
* @param thisRetryCount 当前重试次数
* @return 是否成功
* @return 是否成功
*/
*/
public
static
boolean
retryRpcHost
(
PluginClient
client
,
String
host
,
int
retryTotalCount
,
int
thisRetryCount
)
throws
InterruptedException
{
public
static
boolean
retryRpcHost
(
PluginClient
client
,
String
host
,
String
runKey
,
int
retryTotalCount
,
int
thisRetryCount
)
throws
InterruptedException
{
log
.
info
(
"------当前plugin-rpc的请求的地址为:{}-------"
,
host
);
log
.
info
(
"------当前plugin-rpc的请求的地址为:{}-------"
,
host
);
try
{
try
{
client
.
send
(
host
,
PluginBeat
.
PLUGIN_RPC_REQUEST_PACKET
);
client
.
send
(
host
,
PluginBeat
.
PLUGIN_RPC_REQUEST_PACKET
);
...
@@ -50,9 +58,14 @@ public class ClientRpcUtil {
...
@@ -50,9 +58,14 @@ public class ClientRpcUtil {
}
catch
(
Exception
e
)
{
}
catch
(
Exception
e
)
{
//当前重试次数 小于等于总共的重试次数时
//当前重试次数 小于等于总共的重试次数时
if
(
thisRetryCount
<=
retryTotalCount
)
{
if
(
thisRetryCount
<=
retryTotalCount
)
{
log
.
error
(
"{},plugin-rpc通道建立时出现异常,异常信息为:{},开始第{}次重试!"
,
host
,
RPCLogUtil
.
getMessage
(
e
),
thisRetryCount
);
String
format
=
String
.
format
(
"%s,plugin-rpc通道建立时出现异常,异常信息为:%s,开始第%s次重试!"
,
host
,
RPCLogUtil
.
getMessage
(
e
),
thisRetryCount
);
log
.
error
(
format
);
//redis模板
StringRedisTemplate
stringRedisTemplate
=
RpcSpringUtil
.
getBean
(
StringRedisTemplate
.
class
);
RunLog
runLog
=
RunLog
.
builder
().
runLog
(
format
).
isEnd
(
false
).
build
();
stringRedisTemplate
.
opsForList
().
rightPush
(
runKey
,
JSON
.
toJSONString
(
runLog
,
WriteClassName
));
Thread
.
sleep
(
1000
*
thisRetryCount
);
Thread
.
sleep
(
1000
*
thisRetryCount
);
retryRpcHost
(
client
,
host
,
retryTotalCount
,
++
thisRetryCount
);
retryRpcHost
(
client
,
host
,
runKey
,
retryTotalCount
,
++
thisRetryCount
);
}
}
log
.
error
(
"----与主机【{}】建立通道,总共【{}】次,全部失败,开始挑选下一个负载均衡方案重试,请稍后------"
,
host
,
retryTotalCount
);
log
.
error
(
"----与主机【{}】建立通道,总共【{}】次,全部失败,开始挑选下一个负载均衡方案重试,请稍后------"
,
host
,
retryTotalCount
);
return
false
;
return
false
;
...
...
byit-plugin-core/myth-plugin-rpc-client/src/main/java/com/byit/utils/RpcSpringUtil.java
0 → 100644
View file @
5feb60ca
package
com
.
byit
.
utils
;
import
org.springframework.beans.BeansException
;
import
org.springframework.context.ApplicationContext
;
import
org.springframework.context.ApplicationContextAware
;
import
java.util.Map
;
public
class
RpcSpringUtil
implements
ApplicationContextAware
{
private
static
ApplicationContext
applicationContext
=
null
;
@Override
public
void
setApplicationContext
(
ApplicationContext
applicationContext
)
throws
BeansException
{
if
(
RpcSpringUtil
.
applicationContext
==
null
)
{
RpcSpringUtil
.
applicationContext
=
applicationContext
;
}
}
/**
* 获取applicationContext
* @return 返回applicationContext
*/
public
static
ApplicationContext
getApplicationContext
()
{
return
applicationContext
;
}
/**
* 通过name获取 Bean.
* @param name
* @return
*/
public
static
Object
getBean
(
String
name
){
return
getApplicationContext
().
getBean
(
name
);
}
/**
* 通过class获取Bean.
* @param clazz
* @param <T>
* @return
*/
public
static
<
T
>
T
getBean
(
Class
<
T
>
clazz
){
return
getApplicationContext
().
getBean
(
clazz
);
}
/**
* 通过name,以及Clazz返回指定的Bean
* @param name
* @param clazz
* @param <T>
* @return
*/
public
static
<
T
>
T
getBean
(
String
name
,
Class
<
T
>
clazz
){
return
getApplicationContext
().
getBean
(
name
,
clazz
);
}
/**
* 获取实现某个接口的类
* @param clazz
* @param <T>
* @return
*/
public
static
<
T
>
Map
<
String
,
T
>
getBeansOfType
(
Class
<
T
>
clazz
){
return
getApplicationContext
().
getBeansOfType
(
clazz
);
}
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment