sijie closed pull request #1899: allowing users to specify function jar in yml
file
URL: https://github.com/apache/incubator-pulsar/pull/1899
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
diff --git
a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java
b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java
index 08269567f8..20499e4212 100644
---
a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java
+++
b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java
@@ -40,6 +40,7 @@
import org.apache.pulsar.functions.api.utils.DefaultSerDe;
import org.apache.pulsar.functions.proto.Function.FunctionDetails;
import org.apache.pulsar.functions.utils.Reflections;
+import org.apache.pulsar.functions.utils.Utils;
import org.powermock.api.mockito.PowerMockito;
import org.powermock.core.classloader.annotations.PowerMockIgnore;
import org.powermock.core.classloader.annotations.PrepareForTest;
@@ -76,7 +77,7 @@
* Unit test of {@link CmdFunctions}.
*/
@Slf4j
-@PrepareForTest({ CmdFunctions.class, Reflections.class,
StorageClientBuilder.class })
+@PrepareForTest({ CmdFunctions.class, Reflections.class,
StorageClientBuilder.class, Utils.class})
@PowerMockIgnore({ "javax.management.*", "javax.ws.*",
"org.apache.logging.log4j.*" })
public class CmdFunctionsTest {
@@ -127,6 +128,7 @@ public void setup() throws Exception {
when(Reflections.classImplementsIface(anyString(),
any())).thenReturn(true);
when(Reflections.createInstance(eq(DummyFunction.class.getName()),
any(File.class))).thenReturn(new DummyFunction());
when(Reflections.createInstance(eq(DefaultSerDe.class.getName()),
any(File.class))).thenReturn(new DefaultSerDe(String.class));
+ PowerMockito.stub(PowerMockito.method(Utils.class,
"fileExists")).toReturn(true);
}
// @Test
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
index 30916cdfff..e726cce919 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
@@ -79,6 +79,7 @@
import static org.apache.bookkeeper.common.concurrent.FutureUtils.result;
import static org.apache.pulsar.common.naming.TopicName.DEFAULT_NAMESPACE;
import static org.apache.pulsar.common.naming.TopicName.PUBLIC_TENANT;
+import static org.apache.pulsar.functions.utils.Utils.fileExists;
@Slf4j
@Parameters(commandDescription = "Interface for managing Pulsar Functions
(lightweight, Lambda-style compute processes that work with Pulsar)")
@@ -361,15 +362,18 @@ void processArguments() throws Exception {
functionConfig.setAutoAck(true);
}
-
if (null != jarFile) {
- functionConfig.setRuntime(FunctionConfig.Runtime.JAVA);
- userCodeFile = jarFile;
- } else if (null != pyFile) {
- functionConfig.setRuntime(FunctionConfig.Runtime.PYTHON);
- userCodeFile = pyFile;
- } else {
- throw new ParameterException("Either a Java jar or a Python
file needs to be specified for the function");
+ functionConfig.setJar(jarFile);
+ }
+
+ if (null != pyFile) {
+ functionConfig.setPy(pyFile);
+ }
+
+ if (functionConfig.getJar() != null) {
+ userCodeFile = functionConfig.getJar();
+ } else if (functionConfig.getPy() != null) {
+ userCodeFile = functionConfig.getPy();
}
// infer default vaues
@@ -378,8 +382,22 @@ void processArguments() throws Exception {
protected void validateFunctionConfigs(FunctionConfig functionConfig) {
+ if (functionConfig.getJar() != null && functionConfig.getPy() !=
null) {
+ throw new ParameterException("Either a Java jar or a Python
file needs to"
+ + " be specified for the function. Cannot specify
both.");
+ }
+
+ if (functionConfig.getJar() == null && functionConfig.getPy() ==
null) {
+ throw new ParameterException("Either a Java jar or a Python
file needs to"
+ + " be specified for the function. Please specify
one.");
+ }
+
+ if (!fileExists(userCodeFile)) {
+ throw new ParameterException("File " + userCodeFile + " does
not exist");
+ }
+
if (functionConfig.getRuntime() == FunctionConfig.Runtime.JAVA) {
- File file = new File(jarFile);
+ File file = new File(functionConfig.getJar());
ClassLoader userJarLoader;
try {
userJarLoader = Reflections.loadJar(file);
@@ -394,6 +412,7 @@ protected void validateFunctionConfigs(FunctionConfig
functionConfig) {
// Need to load jar and set context class loader before calling
ConfigValidation.validateConfig(functionConfig,
functionConfig.getRuntime().name());
} catch (Exception e) {
+ log.info("ex: {}", e, e);
throw new ParameterException(e.getMessage());
}
}
@@ -416,6 +435,12 @@ private void inferMissingArguments(FunctionConfig
functionConfig) {
functionConfig.setParallelism(1);
}
+ if (functionConfig.getJar() != null) {
+ functionConfig.setRuntime(FunctionConfig.Runtime.JAVA);
+ } else if (functionConfig.getPy() != null) {
+ functionConfig.setRuntime(FunctionConfig.Runtime.PYTHON);
+ }
+
WindowConfig windowConfig = functionConfig.getWindowConfig();
if (windowConfig != null) {
WindowUtils.inferDefaultConfigs(windowConfig);
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
index fdbe291850..607a4d239e 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
@@ -51,6 +51,7 @@
import static org.apache.pulsar.common.naming.TopicName.DEFAULT_NAMESPACE;
import static org.apache.pulsar.common.naming.TopicName.PUBLIC_TENANT;
import static
org.apache.pulsar.functions.utils.Utils.convertProcessingGuarantee;
+import static org.apache.pulsar.functions.utils.Utils.fileExists;
import static org.apache.pulsar.functions.utils.Utils.getSinkType;
import static org.apache.pulsar.functions.utils.Utils.loadConfig;
@@ -175,7 +176,7 @@ void processArguments() throws Exception {
if (null != tenant) {
sinkConfig.setTenant(tenant);
}
-
+
if (null != namespace) {
sinkConfig.setNamespace(namespace);
}
@@ -208,8 +209,8 @@ void processArguments() throws Exception {
sinkConfig.setParallelism(parallelism);
}
- if (null == jarFile) {
- throw new IllegalArgumentException("Connector JAR not
specfied");
+ if (null != jarFile) {
+ sinkConfig.setJar(jarFile);
}
sinkConfig.setResources(new
org.apache.pulsar.functions.utils.Resources(cpu, ram, disk));
@@ -233,7 +234,15 @@ private void inferMissingArguments(SinkConfig sinkConfig) {
}
protected void validateSinkConfigs(SinkConfig sinkConfig) {
- File file = new File(jarFile);
+ if (null == sinkConfig.getJar()) {
+ throw new ParameterException("Sink jar not specfied");
+ }
+
+ if (!fileExists(sinkConfig.getJar())) {
+ throw new ParameterException("Jar file " + sinkConfig.getJar()
+ " does not exist");
+ }
+
+ File file = new File(sinkConfig.getJar());
ClassLoader userJarLoader;
try {
userJarLoader = Reflections.loadJar(file);
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
index 97d5fd93b4..ffc48917c4 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
@@ -47,6 +47,7 @@
import static org.apache.pulsar.common.naming.TopicName.DEFAULT_NAMESPACE;
import static org.apache.pulsar.common.naming.TopicName.PUBLIC_TENANT;
import static
org.apache.pulsar.functions.utils.Utils.convertProcessingGuarantee;
+import static org.apache.pulsar.functions.utils.Utils.fileExists;
import static org.apache.pulsar.functions.utils.Utils.getSourceType;
import static org.apache.pulsar.functions.utils.Utils.loadConfig;
@@ -191,8 +192,8 @@ void processArguments() throws Exception {
sourceConfig.setParallelism(parallelism);
}
- if (null == jarFile) {
- throw new ParameterException("Source JAR not specfied");
+ if (jarFile != null) {
+ sourceConfig.setJar(jarFile);
}
sourceConfig.setResources(new
org.apache.pulsar.functions.utils.Resources(cpu, ram, disk));
@@ -216,6 +217,14 @@ private void inferMissingArguments(SourceConfig
sourceConfig) {
}
protected void validateSourceConfigs(SourceConfig sourceConfig) {
+ if (null == sourceConfig.getJar()) {
+ throw new ParameterException("Source jar not specfied");
+ }
+
+ if (!fileExists(sourceConfig.getJar())) {
+ throw new ParameterException("Jar file " +
sourceConfig.getJar() + " does not exist");
+ }
+
File file = new File(jarFile);
ClassLoader userJarLoader;
try {
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
index be40e01d05..9671d50cf9 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
@@ -27,6 +27,7 @@
import org.apache.pulsar.functions.api.SerDe;
import org.apache.pulsar.functions.utils.validation.ConfigValidation;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.NotNull;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isFileExists;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isImplementationOfClass;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isImplementationOfClasses;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isListEntryCustom;
@@ -95,4 +96,8 @@
private WindowConfig windowConfig;
@isPositiveNumber
private Long timeoutMs;
+ @isFileExists
+ private String jar;
+ @isFileExists
+ private String py;
}
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SinkConfig.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SinkConfig.java
index 9fa030771e..f4ee77cce4 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SinkConfig.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SinkConfig.java
@@ -24,6 +24,7 @@
import lombok.Setter;
import lombok.ToString;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.NotNull;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isFileExists;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isImplementationOfClass;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isMapEntryCustom;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isPositiveNumber;
@@ -60,4 +61,6 @@
private FunctionConfig.ProcessingGuarantees processingGuarantees;
@isValidResources
private Resources resources;
+ @isFileExists
+ private String jar;
}
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SourceConfig.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SourceConfig.java
index cdede6434e..295f339183 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SourceConfig.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/SourceConfig.java
@@ -25,6 +25,7 @@
import lombok.ToString;
import org.apache.pulsar.functions.api.SerDe;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.NotNull;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isFileExists;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isImplementationOfClass;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isPositiveNumber;
import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidResources;
@@ -62,4 +63,6 @@
private FunctionConfig.ProcessingGuarantees processingGuarantees;
@isValidResources
private Resources resources;
+ @isFileExists
+ private String jar;
}
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/Utils.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/Utils.java
index 3f7c2e47dc..c534ed3db3 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/Utils.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/Utils.java
@@ -205,4 +205,8 @@ public static Runtime convertRuntime(FunctionConfig.Runtime
runtime) {
return typeArg;
}
+
+ public static boolean fileExists(String file) {
+ return new File(file).exists();
+ }
}
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ConfigValidationAnnotations.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ConfigValidationAnnotations.java
index 934f6f55e2..e6e0583a70 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ConfigValidationAnnotations.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ConfigValidationAnnotations.java
@@ -188,6 +188,15 @@
ConfigValidation.Runtime targetRuntime() default
ConfigValidation.Runtime.ALL;
}
+ /**
+ * check if file exists
+ */
+ @Retention(RetentionPolicy.RUNTIME)
+ @Target(ElementType.FIELD)
+ public @interface isFileExists {
+ Class<?> validatorClass() default ValidatorImpls.FileValidator.class;
+ }
+
/**
* checks function config as a whole to make sure all fields are valid
*/
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ValidatorImpls.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ValidatorImpls.java
index 3e0a168196..2ee6044c54 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ValidatorImpls.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/validation/ValidatorImpls.java
@@ -33,12 +33,14 @@
import org.apache.pulsar.functions.utils.Utils;
import org.apache.pulsar.functions.utils.WindowConfig;
+import java.io.File;
import java.lang.reflect.InvocationTargetException;
import java.util.Arrays;
import java.util.Collection;
import java.util.HashSet;
import java.util.Map;
+import static org.apache.pulsar.functions.utils.Utils.fileExists;
import static org.apache.pulsar.functions.utils.Utils.getSinkType;
import static org.apache.pulsar.functions.utils.Utils.getSourceType;
@@ -747,9 +749,22 @@ public void validateField(String name, Object o) {
}
}
+ public static class FileValidator extends Validator {
+ @Override
+ public void validateField(String name, Object o) {
+ if (o == null) {
+ return;
+ }
+ new StringValidator().validateField(name, o);
+ if (!fileExists((String) o)) {
+ throw new IllegalArgumentException
+ (String.format("File %s specified in field '%s' does
not exist", o, name));
+ }
+ }
+ }
- /**
+ /**
* Validates basic types.
*/
public static class SimpleTypeValidator extends Validator {
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services