aggregation-platform/src/com/platform/service/OracleStatusService.java

150 lines
4.0 KiB
Java

package com.platform.service;
import io.fabric8.kubernetes.api.model.Pod;
import io.fabric8.kubernetes.api.model.ReplicationController;
import java.util.Hashtable;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Timer;
import java.util.TimerTask;
import com.platform.entities.OracleConnectorParams;
import com.platform.kubernetes.SimpleKubeClient;
import com.platform.oracle.OracleConnector;
import com.platform.utils.Configs;
public class OracleStatusService {
private static Map<String, Timer> alliveTask = new Hashtable<String, Timer>();
public final static int EXEC_TIME = 10;// 连接多少次后不成功,取消链接
public final static long INTERVAL_TIME = 60 * 1000;// 每隔多少毫秒执行一次连接任务
public final static long DELAY_TIME = 0; // 延迟多少秒后执行
public void connectToOracle(String replicasName) {
SimpleKubeClient sKubeClient = new SimpleKubeClient();
ReplicationController replicationController = sKubeClient
.getReplicationController(replicasName);
if (alliveTask.containsKey(replicasName)) {
killAlliveTask(replicasName);
}
if (null != replicationController) {
List<Pod> filterPods = sKubeClient
.getPodsForApplicaList(replicationController);
if (filterPods != null && filterPods.size() > 0) {
Pod pod = filterPods.get(0);
OracleConnectorParams orp = new OracleConnectorParams(
String.valueOf(sKubeClient.getPodContainerport(pod)),
sKubeClient.getPodHostIp(pod), replicasName);
Timer timer = new Timer();
alliveTask.put(replicasName, timer);
timer.schedule(new connectTask(orp, sKubeClient), DELAY_TIME,
INTERVAL_TIME);
}
}
}
public void cancelToOracle(String replicasName, String operate) {
if (operate.equals("stop")) {
SimpleKubeClient sKubeClient = new SimpleKubeClient();
sKubeClient.updateOrAddReplicasLabelById(replicasName, "status",
"0");
}
killAlliveTask(replicasName);
}
/**
* 取消并移除指定定时任务
*
*
* @param taskName
*/
public void killAlliveTask(String taskName) {
if (alliveTask.containsKey(taskName)) {
alliveTask.get(taskName).cancel();
alliveTask.remove(taskName);
}
}
public void killAlliveTasks(String... tasksName) {
for (String taskName : tasksName)
killAlliveTask(taskName);
}
/**
* 清空定时任务
*/
public void cleanUpAlliveTask() {
Iterator<Map.Entry<String, Timer>> iterator = alliveTask.entrySet()
.iterator();
while (iterator.hasNext()) {
Map.Entry<String, Timer> entry = iterator.next();
entry.getValue().cancel();
}
alliveTask.clear();
}
/**
* 链接oracle任务类
*
* @author wuming
*
*/
class connectTask extends TimerTask {
private String taskName;
private int count;
private OracleConnectorParams ocp;
private SimpleKubeClient client;
public connectTask(OracleConnectorParams ocp, SimpleKubeClient client) {
this.taskName = ocp.getName();
this.ocp = ocp;
this.count = 0;
this.client = client;
}
@Override
public void run() {
if (count == EXEC_TIME && alliveTask.containsKey(taskName)) {
killAlliveTask(taskName);
client.updateOrAddReplicasLabelById(taskName, "status", "1");
} else {
String url = "jdbc:oracle:thin:@" + ocp.getIp() + ":"
+ ocp.getPort() + ":" + ocp.getDatabase();
boolean flag = OracleConnector.canConnect(url, ocp.getUser(),
ocp.getPassword());
String message = "失败";
if (flag && alliveTask.containsKey(taskName)) {
client.updateOrAddReplicasLabelById(taskName, "status", "2");
message = "成功";
killAlliveTask(taskName); // 连接成功,取消连接
}
Configs.CONSOLE_LOGGER.info("连接到数据库服务: " + taskName
+ "\t[连接结果: " + message + "]");
}
count++;
}
@Override
public boolean cancel() {
System.out.println("aaaaaaa");
if (client != null)
this.client.close();
return super.cancel();
}
public String getTaskName() {
return taskName;
}
public void setTaskName(String taskName) {
this.taskName = taskName;
}
public int getCount() {
return count;
}
}
}