Skip to content
项目
群组
代码片段
帮助
正在加载...
登录
切换导航
X
XXL-JOB
项目
项目
详情
活动
周期分析
仓库
仓库
文件
提交
分支
标签
贡献者
分枝图
比较
统计图
议题
0
议题
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
CI / CD
CI / CD
流水线
作业
日程
统计图
Wiki
Wiki
代码片段
代码片段
成员
成员
折叠边栏
关闭边栏
活动
分枝图
统计图
创建新议题
作业
提交
议题看板
打开侧边栏
靳帅
XXL-JOB
Commits
dbbb1f06
提交
dbbb1f06
authored
1月 16, 2016
作者:
xueli.xue
浏览文件
操作
浏览文件
下载
差异文件
Merge branch 'master' of
https://github.com/xuxueli/xxl-job.git
上级
554429ea
ed226c60
隐藏空白字符变更
内嵌
并排
正在显示
3 个修改的文件
包含
194 行增加
和
0 行删除
+194
-0
LocalJobBean.java
...n/src/main/java/com/xxl/job/service/job/LocalJobBean.java
+51
-0
LocalJobBeanB.java
.../src/main/java/com/xxl/job/service/job/LocalJobBeanB.java
+49
-0
HandlerThread.java
...c/main/java/com/xxl/job/client/handler/HandlerThread.java
+94
-0
没有找到文件。
xxl-job-admin/src/main/java/com/xxl/job/service/job/LocalJobBean.java
0 → 100644
浏览文件 @
dbbb1f06
package
com
.
xxl
.
job
.
service
.
job
;
import
java.util.HashMap
;
import
java.util.Map
;
import
java.util.Map.Entry
;
import
java.util.concurrent.TimeUnit
;
import
org.quartz.DisallowConcurrentExecution
;
import
org.quartz.JobExecutionContext
;
import
org.quartz.JobExecutionException
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.scheduling.quartz.QuartzJobBean
;
/**
* http job bean
* @author xuxueli 2015-12-17 18:20:34
*/
@DisallowConcurrentExecution
// 串行;线程数要多配置几个,否则不生效;
public
class
LocalJobBean
extends
QuartzJobBean
{
private
static
Logger
logger
=
LoggerFactory
.
getLogger
(
LocalJobBean
.
class
);
@Override
protected
void
executeInternal
(
JobExecutionContext
context
)
throws
JobExecutionException
{
String
triggerKey
=
context
.
getTrigger
().
getKey
().
getName
();
String
triggerGroup
=
context
.
getTrigger
().
getKey
().
getGroup
();
Map
<
String
,
Object
>
jobDataMap
=
context
.
getMergedJobDataMap
().
getWrappedMap
();
// jobDataMap 2 params
Map
<
String
,
String
>
params
=
new
HashMap
<
String
,
String
>();
if
(
jobDataMap
!=
null
&&
jobDataMap
.
size
()>
0
)
{
for
(
Entry
<
String
,
Object
>
item
:
jobDataMap
.
entrySet
())
{
params
.
put
(
item
.
getKey
(),
String
.
valueOf
(
item
.
getValue
()));
}
}
try
{
TimeUnit
.
SECONDS
.
sleep
(
5
);
}
catch
(
InterruptedException
e
)
{
e
.
printStackTrace
();
}
logger
.
info
(
">>>>>>>>>>> xxl-job run :jobId:{}, group:{}, jobDataMap:{}"
,
new
Object
[]{
triggerKey
,
triggerGroup
,
jobDataMap
});
}
}
\ No newline at end of file
xxl-job-admin/src/main/java/com/xxl/job/service/job/LocalJobBeanB.java
0 → 100644
浏览文件 @
dbbb1f06
package
com
.
xxl
.
job
.
service
.
job
;
import
java.util.HashMap
;
import
java.util.Map
;
import
java.util.Map.Entry
;
import
java.util.concurrent.TimeUnit
;
import
org.quartz.JobExecutionContext
;
import
org.quartz.JobExecutionException
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
org.springframework.scheduling.quartz.QuartzJobBean
;
/**
* http job bean
* @author xuxueli 2015-12-17 18:20:34
*/
public
class
LocalJobBeanB
extends
QuartzJobBean
{
private
static
Logger
logger
=
LoggerFactory
.
getLogger
(
LocalJobBeanB
.
class
);
@Override
protected
void
executeInternal
(
JobExecutionContext
context
)
throws
JobExecutionException
{
String
triggerKey
=
context
.
getTrigger
().
getKey
().
getName
();
String
triggerGroup
=
context
.
getTrigger
().
getKey
().
getGroup
();
Map
<
String
,
Object
>
jobDataMap
=
context
.
getMergedJobDataMap
().
getWrappedMap
();
// jobDataMap 2 params
Map
<
String
,
String
>
params
=
new
HashMap
<
String
,
String
>();
if
(
jobDataMap
!=
null
&&
jobDataMap
.
size
()>
0
)
{
for
(
Entry
<
String
,
Object
>
item
:
jobDataMap
.
entrySet
())
{
params
.
put
(
item
.
getKey
(),
String
.
valueOf
(
item
.
getValue
()));
}
}
try
{
TimeUnit
.
SECONDS
.
sleep
(
5
);
}
catch
(
InterruptedException
e
)
{
e
.
printStackTrace
();
}
logger
.
info
(
">>>>>>>>>>> xxl-job run :jobId:{}, group:{}, jobDataMap:{}"
,
new
Object
[]{
triggerKey
,
triggerGroup
,
jobDataMap
});
}
}
\ No newline at end of file
xxl-job-client/src/main/java/com/xxl/job/client/handler/HandlerThread.java
0 → 100644
浏览文件 @
dbbb1f06
package
com
.
xxl
.
job
.
client
.
handler
;
import
java.io.PrintWriter
;
import
java.io.StringWriter
;
import
java.util.HashMap
;
import
java.util.Map
;
import
java.util.concurrent.LinkedBlockingQueue
;
import
java.util.concurrent.TimeUnit
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
import
com.xxl.job.client.handler.IJobHandler.JobHandleStatus
;
import
com.xxl.job.client.util.HttpUtil
;
/**
* handler thread
* @author xuxueli 2016-1-16 19:52:47
*/
public
class
HandlerThread
extends
Thread
{
private
static
Logger
logger
=
LoggerFactory
.
getLogger
(
HandlerThread
.
class
);
private
IJobHandler
handler
;
private
LinkedBlockingQueue
<
Map
<
String
,
String
>>
handlerDataQueue
;
public
HandlerThread
(
IJobHandler
handler
)
{
this
.
handler
=
handler
;
handlerDataQueue
=
new
LinkedBlockingQueue
<
Map
<
String
,
String
>>();
}
public
void
pushData
(
Map
<
String
,
String
>
param
)
{
handlerDataQueue
.
offer
(
param
);
}
int
i
=
1
;
@Override
public
void
run
()
{
try
{
i
++;
Map
<
String
,
String
>
handlerData
=
handlerDataQueue
.
poll
();
if
(
handlerData
!=
null
)
{
String
trigger_log_url
=
handlerData
.
get
(
HandlerRepository
.
TRIGGER_LOG_URL
);
String
trigger_log_id
=
handlerData
.
get
(
HandlerRepository
.
TRIGGER_LOG_ID
);
String
handler_params
=
handlerData
.
get
(
HandlerRepository
.
HANDLER_PARAMS
);
// parse param
String
[]
handlerParams
=
null
;
if
(
handler_params
!=
null
&&
handler_params
.
trim
().
length
()>
0
)
{
handlerParams
=
handler_params
.
split
(
","
);
}
else
{
handlerParams
=
new
String
[
0
];
}
// handle job
JobHandleStatus
_status
=
JobHandleStatus
.
FAIL
;
String
_msg
=
null
;
try
{
_status
=
handler
.
handle
(
handlerParams
);
}
catch
(
Exception
e
)
{
logger
.
info
(
"HandlerThread Exception:"
,
e
);
StringWriter
out
=
new
StringWriter
();
e
.
printStackTrace
(
new
PrintWriter
(
out
));
_msg
=
out
.
toString
();
}
// callback handler info
String
callback_response
[]
=
null
;
try
{
HashMap
<
String
,
String
>
params
=
new
HashMap
<
String
,
String
>();
params
.
put
(
HandlerRepository
.
TRIGGER_LOG_ID
,
trigger_log_id
);
params
.
put
(
HttpUtil
.
status
,
_status
.
name
());
params
.
put
(
HttpUtil
.
msg
,
_msg
);
callback_response
=
HttpUtil
.
post
(
trigger_log_url
,
params
);
}
catch
(
Exception
e
)
{
logger
.
info
(
"HandlerThread Exception:"
,
e
);
}
logger
.
info
(
"<<<<<<<<<<< xxl-job thread handle, handlerData:{}, callback_status:{}, callback_msg:{}, callback_response:{}, thread:{}"
,
new
Object
[]{
handlerData
,
_status
,
_msg
,
callback_response
,
this
});
}
else
{
try
{
TimeUnit
.
MILLISECONDS
.
sleep
(
i
*
100
);
}
catch
(
InterruptedException
e
)
{
e
.
printStackTrace
();
}
if
(
i
>
5
)
{
i
=
0
;
}
}
}
catch
(
Exception
e
)
{
logger
.
info
(
"HandlerThread Exception:"
,
e
);
}
}
}
编写
预览
Markdown
格式
0%
重试
或
添加新文件
添加附件
取消
您添加了
0
人
到此讨论。请谨慎行事。
请先完成此评论的编辑!
取消
请
注册
或者
登录
后发表评论