* lib-scheduler-engines/src/main/java/com/fourelementscapital/scheduler/engines/PythonScript.java

* lib-scheduler-engines/src/main/java/com/fourelementscapital/scheduler/group/PythonScriptTask.java
      Added new: the new python engine based on cpython

  * lib-scheduler-queue/src/main/java/com/fourelementscapital/scheduler/peer/QueueFactory.java
      Added: the new factory for cpython

  * lib-scheduler-queue/src/main/java/com/fourelementscapital/scheduler/ScheduledTaskFactory.java
      Adjusted: accordingly to hook up the new engine
            Associated SQL statements needed in bbsync database:

              update scheduler_group set enginetype = 'pscript4pythonengine'
                where taskuid = 'rscript4rserveunix_python';

              insert into scheduler_taskpeers( taskuid, peername )
                values( 'rscript4rserveunix_python', '4ecapsvid13' );
This commit is contained in:
2021-12-16 15:31:45 +08:00
parent 15c5bd748d
commit adf29ba998
4 changed files with 279 additions and 41 deletions
@@ -20,6 +20,7 @@ import com.fourelementscapital.db.SchedulerDB;
import com.fourelementscapital.scheduler.engines.ScheduledTask;
import com.fourelementscapital.scheduler.group.BBDownloadScheduledTask;
import com.fourelementscapital.scheduler.group.DirectRServeExecuteRUnix;
import com.fourelementscapital.scheduler.group.PythonScriptTask;
import com.fourelementscapital.scheduler.group.REngineScriptTask;
import com.fourelementscapital.scheduler.group.RScriptScheduledTask;
import com.fourelementscapital.scheduler.group.RServeScheduledTask;
@@ -38,7 +39,8 @@ public class ScheduledTaskFactory {
init();
}
public synchronized void refreshTaskLoaded(){
//public synchronized void refreshTaskLoaded(){
public void refreshTaskLoaded(){
synchronized(scheduledTasks){
if(allTasks.size()>0 || scheduledTasks.size()>0){
allTasks.clear();
@@ -58,37 +60,38 @@ public class ScheduledTaskFactory {
Vector venabled=SchedulerEngine.getEnabledTaskTypes();
SchedulerDB sdb=SchedulerDB.getSchedulerDB();
try{
sdb.connectDB();
Vector allgroups=sdb.getAllGroups();
for(Iterator i=allgroups.iterator();i.hasNext();){
Map data=(Map)i.next();
String taskuid=(String)data.get("taskuid");
String name=(String)data.get("name");
String enginetype=(String)data.get("enginetype");
ScheduledTask st=null;
if(enginetype.equalsIgnoreCase("rscript")){st=new RScriptScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("rscript4rserve")){st=new RServeScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("bb_download")){st=new BBDownloadScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase(RServeUnixTask.ENGINE_NAME)){st=new RServeUnixTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("direct_script")){st=new REngineScriptTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("direct_script_unix")){st=new DirectRServeExecuteRUnix(name,taskuid);}
//if(enginetype.equalsIgnoreCase("rscript4rconsole")){st=new RConsoleScript(name,taskuid);}
if(st!=null){
allTasks.add(st);
}
}
}catch(Exception e){
}finally{
try{
sdb.closeDB();
}catch(Exception e){}
}
SchedulerDB sdb=SchedulerDB.getSchedulerDB();
try{
sdb.connectDB();
Vector allgroups=sdb.getAllGroups();
for(Iterator i=allgroups.iterator();i.hasNext();){
Map data=(Map)i.next();
String taskuid=(String)data.get("taskuid");
String name=(String)data.get("name");
String enginetype=(String)data.get("enginetype");
ScheduledTask st=null;
if(enginetype.equalsIgnoreCase("rscript")){st=new RScriptScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("rscript4rserve")){st=new RServeScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("bb_download")){st=new BBDownloadScheduledTask(name,taskuid); }
if(enginetype.equalsIgnoreCase(RServeUnixTask.ENGINE_NAME)){st=new RServeUnixTask(name,taskuid); } //rscript4rserveunix
if(enginetype.equalsIgnoreCase("direct_script")){st=new REngineScriptTask(name,taskuid); }
if(enginetype.equalsIgnoreCase("direct_script_unix")){st=new DirectRServeExecuteRUnix(name,taskuid);}
if(enginetype.equalsIgnoreCase(PythonScriptTask.ENGINE_NAME)){st=new PythonScriptTask(name,taskuid);}
//if(enginetype.equalsIgnoreCase("rscript4rconsole")){st=new RConsoleScript(name,taskuid);}
if(st!=null){
allTasks.add(st);
}
}
}catch(Exception e){
e.printStackTrace();
}finally{
try{
sdb.closeDB();
}catch(Exception e){}
}
@@ -15,6 +15,9 @@ import java.util.Map;
import java.util.TreeMap;
import java.util.Vector;
import java.util.Set;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@@ -25,6 +28,7 @@ import com.fourelementscapital.scheduler.engines.ScheduledTask;
import com.fourelementscapital.scheduler.engines.StackFrame;
import com.fourelementscapital.scheduler.p2p.P2PService;
import com.fourelementscapital.scheduler.p2p.peer.PeerSpecificConfigurations;
import com.fourelementscapital.scheduler.exception.ExceptionPeerRejected;
public class QueueFactory {
@@ -43,16 +47,18 @@ public class QueueFactory {
public QueueFactory(){
if(queueHandlers.size()==0){
initQueue();
initQueue();
}
}
public QueueAbstract getQueue(String taskuid){
QueueAbstract queueAbs = (QueueAbstract)queueHandlers.get(taskuid);
return (QueueAbstract)queueHandlers.get(taskuid);
}
public TreeMap getQueue(){
return queueHandlers;
}
@@ -144,17 +150,16 @@ public class QueueFactory {
}
private void initQueue(){
SchedulerDB sdb=SchedulerDB.getSchedulerDB();
try{
sdb.connectDB();
sdb.removeAllPeerThreadStatus(P2PService.getComputerName());
}catch(Exception e){ }finally{
try{sdb.closeDB();}catch(Exception e1){}
}
//over-rides
QueueAbstract rserv=new QueueAbstract("rserv"){
QueueAbstract rserv=new QueueAbstract("rserv"){ /////////////////////////////////////////////////////////////////////////
public Vector getTaskUids(){
Vector v=new Vector();
@@ -203,7 +208,6 @@ public class QueueFactory {
};
add2Qhandler(rserv);
//over-rides
QueueAbstract rservunix=new QueueAbstract("rservunix"){
public Vector getTaskUids(){
@@ -245,8 +249,45 @@ public class QueueFactory {
}
};
add2Qhandler(rservunix);
//over-rides
QueueAbstract cpythonlinux = new QueueAbstract("cpythonlinux") {
@Override
public Vector getTaskUids() {
Vector result = new Vector();
SchedulerDB sdb = SchedulerDB.getSchedulerDB();
try {
sdb.connectDB();
Vector< Map<String, String> > grp = sdb.getGroups( "pscript4pythonengine" );
for (Map<String, String> data : grp) {
result.add( data.get( "taskuid" ) );
}
}
catch (Exception e) {
}
finally {
try {
sdb.closeDB();
}
catch (Exception e1) {
}
}
return result;
}
@Override
public int getConcurrentThreads() {
int result = 15; // For now
return result;
}
};
add2Qhandler(cpythonlinux);
//over-rides
QueueAbstract rscript=new QueueAbstract("rscript"){
public Vector getTaskUids(){
Vector v=new Vector();
@@ -278,9 +319,8 @@ public class QueueFactory {
return 1;
}
};
add2Qhandler(rscript);
QueueAbstract othercripts=new QueueAbstract("othercripts"){
public Vector getTaskUids(){
Vector v=new Vector();
@@ -297,12 +337,13 @@ public class QueueFactory {
return 1;
}
};
add2Qhandler(othercripts);
add2Qhandler(othercripts);
}
private void add2Qhandler(QueueAbstract qa){
for(Iterator i=qa.getTaskUids().iterator();i.hasNext();){
queueHandlers.put(i.next(), qa);
}