[Improvement] Merge redundant codes (#7550)
This commit is contained in:
parent
4c918f6167
commit
6ae46f2c1b
|
|
@ -17,6 +17,7 @@
|
|||
|
||||
package org.apache.dolphinscheduler.plugin.task.api;
|
||||
|
||||
import org.apache.dolphinscheduler.spi.task.ResourceInfo;
|
||||
import org.apache.dolphinscheduler.spi.task.request.TaskRequest;
|
||||
|
||||
/**
|
||||
|
|
@ -80,4 +81,21 @@ public abstract class AbstractYarnTask extends AbstractTaskExecutor {
|
|||
* set main jar name
|
||||
*/
|
||||
protected abstract void setMainJarName();
|
||||
|
||||
/**
|
||||
* Get name of jar resource.
|
||||
*
|
||||
* @param mainJar
|
||||
* @return
|
||||
*/
|
||||
protected String getResourceNameOfMainJar(ResourceInfo mainJar) {
|
||||
if (null == mainJar) {
|
||||
throw new RuntimeException("The jar for the task is required.");
|
||||
}
|
||||
|
||||
return mainJar.getId() == 0
|
||||
? mainJar.getRes()
|
||||
// when update resource maybe has error
|
||||
: mainJar.getResourceName().replaceFirst("/", "");
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -116,17 +116,9 @@ public class FlinkTask extends AbstractYarnTask {
|
|||
protected void setMainJarName() {
|
||||
// main jar
|
||||
ResourceInfo mainJar = flinkParameters.getMainJar();
|
||||
if (mainJar != null) {
|
||||
int resourceId = mainJar.getId();
|
||||
String resourceName;
|
||||
if (resourceId == 0) {
|
||||
resourceName = mainJar.getRes();
|
||||
} else {
|
||||
resourceName = mainJar.getResourceName().replaceFirst("/", "");
|
||||
}
|
||||
mainJar.setRes(resourceName);
|
||||
flinkParameters.setMainJar(mainJar);
|
||||
}
|
||||
String resourceName = getResourceNameOfMainJar(mainJar);
|
||||
mainJar.setRes(resourceName);
|
||||
flinkParameters.setMainJar(mainJar);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -119,17 +119,9 @@ public class MapReduceTask extends AbstractYarnTask {
|
|||
protected void setMainJarName() {
|
||||
// main jar
|
||||
ResourceInfo mainJar = mapreduceParameters.getMainJar();
|
||||
if (mainJar != null) {
|
||||
int resourceId = mainJar.getId();
|
||||
String resourceName;
|
||||
if (resourceId == 0) {
|
||||
resourceName = mainJar.getRes();
|
||||
} else {
|
||||
resourceName = mainJar.getResourceName().replaceFirst("/", "");
|
||||
}
|
||||
mainJar.setRes(resourceName);
|
||||
mapreduceParameters.setMainJar(mainJar);
|
||||
}
|
||||
String resourceName = getResourceNameOfMainJar(mainJar);
|
||||
mainJar.setRes(resourceName);
|
||||
mapreduceParameters.setMainJar(mainJar);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
|
|
@ -121,22 +121,9 @@ public class SparkTask extends AbstractYarnTask {
|
|||
protected void setMainJarName() {
|
||||
// main jar
|
||||
ResourceInfo mainJar = sparkParameters.getMainJar();
|
||||
|
||||
if (null == mainJar) {
|
||||
throw new RuntimeException("Spark task jar params is null");
|
||||
}
|
||||
|
||||
int resourceId = mainJar.getId();
|
||||
String resourceName;
|
||||
if (resourceId == 0) {
|
||||
resourceName = mainJar.getRes();
|
||||
} else {
|
||||
//when update resource maybe has error
|
||||
resourceName = mainJar.getResourceName().replaceFirst("/", "");
|
||||
}
|
||||
String resourceName = getResourceNameOfMainJar(mainJar);
|
||||
mainJar.setRes(resourceName);
|
||||
sparkParameters.setMainJar(mainJar);
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
|
|||
Loading…
Reference in New Issue