sijie closed pull request #1897: improving source and sink validation
URL: https://github.com/apache/incubator-pulsar/pull/1897
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/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 b3e4f6bb21..ac4790c5c4 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
@@ -22,8 +22,6 @@
import com.beust.jcommander.ParameterException;
import com.beust.jcommander.Parameters;
import com.beust.jcommander.converters.StringConverter;
-import com.fasterxml.jackson.databind.ObjectMapper;
-import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import com.google.common.annotations.VisibleForTesting;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
@@ -46,7 +44,6 @@
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.functions.instance.InstanceConfig;
import org.apache.pulsar.functions.proto.Function.FunctionDetails;
-import org.apache.pulsar.functions.proto.Function.ProcessingGuarantees;
import org.apache.pulsar.functions.proto.Function.Resources;
import org.apache.pulsar.functions.proto.Function.SinkSpec;
import org.apache.pulsar.functions.proto.Function.SourceSpec;
@@ -259,7 +256,7 @@ void processArguments() throws Exception {
// Initialize config builder either from a supplied YAML config
file or from scratch
if (null != fnConfigFile) {
- functionConfig = loadConfig(new File(fnConfigFile));
+ functionConfig = Utils.loadConfig(fnConfigFile,
FunctionConfig.class);
} else {
functionConfig = new FunctionConfig();
}
@@ -471,6 +468,9 @@ private String getUniqueInput(FunctionConfig
functionConfig) {
protected FunctionDetails convert(FunctionConfig functionConfig)
throws IOException {
+ // check if configs are valid
+ validateFunctionConfigs(functionConfig);
+
Class<?>[] typeArgs = null;
if (functionConfig.getRuntime() == FunctionConfig.Runtime.JAVA) {
// Assuming any external jars are already loaded
@@ -535,11 +535,11 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
functionDetailsBuilder.setLogTopic(functionConfig.getLogTopic());
}
if (functionConfig.getRuntime() != null) {
-
functionDetailsBuilder.setRuntime(convertRuntime(functionConfig.getRuntime()));
+
functionDetailsBuilder.setRuntime(Utils.convertRuntime(functionConfig.getRuntime()));
}
if (functionConfig.getProcessingGuarantees() != null) {
functionDetailsBuilder.setProcessingGuarantees(
-
convertProcessingGuarantee(functionConfig.getProcessingGuarantees()));
+
Utils.convertProcessingGuarantee(functionConfig.getProcessingGuarantees()));
}
Map<String, Object> configs = new HashMap<>();
@@ -600,8 +600,6 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
@Override
void runCmd() throws Exception {
- // check if function configs are valid
- validateFunctionConfigs(functionConfig);
CmdFunctions.startLocalRun(convertProto2(functionConfig),
functionConfig.getParallelism(), brokerServiceUrl,
userCodeFile, admin);
}
@@ -611,8 +609,6 @@ void runCmd() throws Exception {
class CreateFunction extends FunctionDetailsCommand {
@Override
void runCmd() throws Exception {
- // check if function configs are valid
- validateFunctionConfigs(functionConfig);
admin.functions().createFunction(convert(functionConfig),
userCodeFile);
print("Created successfully");
}
@@ -651,8 +647,6 @@ void runCmd() throws Exception {
class UpdateFunction extends FunctionDetailsCommand {
@Override
void runCmd() throws Exception {
- // check if function configs are valid
- validateFunctionConfigs(functionConfig);
admin.functions().updateFunction(convert(functionConfig),
userCodeFile);
print("Updated successfully");
}
@@ -860,30 +854,6 @@ DownloadFunction getDownloader() {
return downloader;
}
- private static FunctionConfig loadConfig(File file) throws IOException {
- ObjectMapper mapper = new ObjectMapper(new YAMLFactory());
- return mapper.readValue(file, FunctionConfig.class);
- }
-
- private static FunctionDetails.Runtime
convertRuntime(FunctionConfig.Runtime runtime) {
- for (FunctionDetails.Runtime type : FunctionDetails.Runtime.values()) {
- if (type.name().equals(runtime.name())) {
- return type;
- }
- }
- throw new RuntimeException("Unrecognized runtime: " + runtime.name());
- }
-
- private static ProcessingGuarantees convertProcessingGuarantee(
- FunctionConfig.ProcessingGuarantees processingGuarantees) {
- for (ProcessingGuarantees type : ProcessingGuarantees.values()) {
- if (type.name().equals(processingGuarantees.name())) {
- return type;
- }
- }
- throw new RuntimeException("Unrecognized processing guarantee: " +
processingGuarantees.name());
- }
-
private void parseFullyQualifiedFunctionName(String fqfn, FunctionConfig
functionConfig) {
String[] args = fqfn.split("/");
if (args.length != 3) {
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 2579c41cb0..fdbe291850 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
@@ -19,32 +19,17 @@
package org.apache.pulsar.admin.cli;
import com.beust.jcommander.Parameter;
+import com.beust.jcommander.ParameterException;
import com.beust.jcommander.Parameters;
import com.beust.jcommander.converters.StringConverter;
-import com.fasterxml.jackson.databind.ObjectMapper;
-import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import com.google.gson.Gson;
import com.google.gson.reflect.TypeToken;
-
-import java.io.File;
-import java.io.IOException;
-import java.lang.reflect.Type;
-import java.net.MalformedURLException;
-import java.util.Arrays;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.function.Consumer;
-
import lombok.Getter;
-
import org.apache.pulsar.client.admin.PulsarAdmin;
import org.apache.pulsar.client.admin.internal.FunctionsImpl;
-import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.functions.api.utils.IdentityFunction;
import org.apache.pulsar.functions.proto.Function;
import org.apache.pulsar.functions.proto.Function.FunctionDetails;
-import org.apache.pulsar.functions.proto.Function.ProcessingGuarantees;
import org.apache.pulsar.functions.proto.Function.Resources;
import org.apache.pulsar.functions.proto.Function.SinkSpec;
import org.apache.pulsar.functions.proto.Function.SourceSpec;
@@ -52,12 +37,22 @@
import org.apache.pulsar.functions.utils.Reflections;
import org.apache.pulsar.functions.utils.SinkConfig;
import org.apache.pulsar.functions.utils.Utils;
-import org.apache.pulsar.io.core.Sink;
+import org.apache.pulsar.functions.utils.validation.ConfigValidation;
-import net.jodah.typetools.TypeResolver;
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Type;
+import java.net.MalformedURLException;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
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.getSinkType;
+import static org.apache.pulsar.functions.utils.Utils.loadConfig;
@Getter
@Parameters(commandDescription = "Interface for managing Pulsar Sinks (Egress
data from Pulsar)")
@@ -115,9 +110,6 @@ void runCmd() throws Exception {
class CreateSink extends SinkCommand {
@Override
void runCmd() throws Exception {
- if (!areAllRequiredFieldsPresentForSink(sinkConfig)) {
- throw new RuntimeException("Missing arguments");
- }
admin.functions().createFunction(createSinkConfig(sinkConfig),
jarFile);
print("Created successfully");
}
@@ -127,9 +119,6 @@ void runCmd() throws Exception {
class UpdateSink extends SinkCommand {
@Override
void runCmd() throws Exception {
- if (!areAllRequiredFieldsPresentForSink(sinkConfig)) {
- throw new RuntimeException("Missing arguments");
- }
admin.functions().updateFunction(createSinkConfig(sinkConfig),
jarFile);
print("Updated successfully");
}
@@ -152,7 +141,7 @@ void runCmd() throws Exception {
@Parameter(names = "--processingGuarantees", description = "The
processing guarantees (aka delivery semantics) applied to the Sink")
protected FunctionConfig.ProcessingGuarantees processingGuarantees;
@Parameter(names = "--parallelism", description = "The sink's
parallelism factor (i.e. the number of sink instances to run)")
- protected String parallelism;
+ protected Integer parallelism;
@Parameter(
names = "--jar",
description = "Path to the jar file for the sink",
@@ -178,21 +167,19 @@ void processArguments() throws Exception {
super.processArguments();
if (null != sinkConfigFile) {
- this.sinkConfig = loadSinkConfig(sinkConfigFile);
+ this.sinkConfig = loadConfig(sinkConfigFile, SinkConfig.class);
} else {
this.sinkConfig = new SinkConfig();
}
if (null != tenant) {
sinkConfig.setTenant(tenant);
- } else if (sinkConfig.getTenant() == null) {
- sinkConfig.setTenant(PUBLIC_TENANT);
}
+
if (null != namespace) {
sinkConfig.setNamespace(namespace);
- } else if (sinkConfig.getNamespace() == null) {
- sinkConfig.setNamespace(DEFAULT_NAMESPACE);
}
+
if (null != name) {
sinkConfig.setName(name);
}
@@ -206,44 +193,25 @@ void processArguments() throws Exception {
Map<String, String> topicsToSerDeClassName = new HashMap<>();
if (null != inputs) {
List<String> inputTopics = Arrays.asList(inputs.split(","));
- inputTopics.forEach(new Consumer<String>() {
- @Override
- public void accept(String s) {
- CmdSinks.validateTopicName(s);
- topicsToSerDeClassName.put(s, "");
- }
- });
+ inputTopics.forEach(s -> topicsToSerDeClassName.put(s, ""));
}
if (null != customSerdeInputString) {
Type type = new TypeToken<Map<String, String>>(){}.getType();
Map<String, String> customSerdeInputMap = new
Gson().fromJson(customSerdeInputString, type);
customSerdeInputMap.forEach((topic, serde) -> {
- CmdSinks.validateTopicName(topic);
topicsToSerDeClassName.put(topic, serde);
});
}
sinkConfig.setTopicToSerdeClassName(topicsToSerDeClassName);
- if (parallelism == null) {
- if (sinkConfig.getParallelism() == 0) {
- sinkConfig.setParallelism(1);
- }
- } else {
- int num = Integer.parseInt(parallelism);
- if (num <= 0) {
- throw new IllegalArgumentException("The parallelism factor
(the number of instances) for the "
- + "connector must be positive");
- }
- sinkConfig.setParallelism(num);
+ if (parallelism != null) {
+ sinkConfig.setParallelism(parallelism);
}
if (null == jarFile) {
throw new IllegalArgumentException("Connector JAR not
specfied");
}
- com.google.common.base.Preconditions.checkArgument(cpu == null ||
cpu > 0, "The cpu allocation for the sink must be positive");
- com.google.common.base.Preconditions.checkArgument(ram == null ||
ram > 0, "The ram allocation for the sink must be positive");
- com.google.common.base.Preconditions.checkArgument(disk == null ||
disk > 0, "The disk allocation for the sink must be positive");
sinkConfig.setResources(new
org.apache.pulsar.functions.utils.Resources(cpu, ram, disk));
if (null != sinkConfigString) {
@@ -251,29 +219,39 @@ public void accept(String s) {
Map<String, Object> sinkConfigMap = new
Gson().fromJson(sinkConfigString, type);
sinkConfig.setConfigs(sinkConfigMap);
}
+
+ inferMissingArguments(sinkConfig);
}
- private Class<?> getSinkType(File file) {
- if (!Reflections.classExistsInJar(file,
sinkConfig.getClassName())) {
- throw new IllegalArgumentException(String.format("Pulsar sink
class %s does not exist in jar %s",
- sinkConfig.getClassName(), jarFile));
- } else if (!Reflections.classInJarImplementsIface(file,
sinkConfig.getClassName(), Sink.class)) {
- throw new IllegalArgumentException(String.format("The Pulsar
sink class %s in jar %s implements " + Sink.class.getName(),
- sinkConfig.getClassName(), jarFile));
+ private void inferMissingArguments(SinkConfig sinkConfig) {
+ if (sinkConfig.getTenant() == null) {
+ sinkConfig.setTenant(PUBLIC_TENANT);
+ }
+ if (sinkConfig.getNamespace() == null) {
+ sinkConfig.setNamespace(DEFAULT_NAMESPACE);
}
+ }
- Object userClass =
Reflections.createInstance(sinkConfig.getClassName(), file);
- Class<?> typeArg;
- Sink sink = (Sink) userClass;
- if (sink == null) {
- throw new IllegalArgumentException(String.format("The Pulsar
sink class %s could not be instantiated from jar %s",
- sinkConfig.getClassName(), jarFile));
+ protected void validateSinkConfigs(SinkConfig sinkConfig) {
+ File file = new File(jarFile);
+ ClassLoader userJarLoader;
+ try {
+ userJarLoader = Reflections.loadJar(file);
+ } catch (MalformedURLException e) {
+ throw new ParameterException("Failed to load user jar " + file
+ " with error " + e.getMessage());
}
- typeArg = TypeResolver.resolveRawArgument(Sink.class,
sink.getClass());
+ // make sure the function class loader is accessible thread-locally
+ Thread.currentThread().setContextClassLoader(userJarLoader);
- return typeArg;
+ try {
+ // Need to load jar and set context class loader before calling
+ ConfigValidation.validateConfig(sinkConfig,
FunctionConfig.Runtime.JAVA.name());
+ } catch (Exception e) {
+ throw new ParameterException(e.getMessage());
+ }
}
+
protected org.apache.pulsar.functions.proto.Function.FunctionDetails
createSinkConfigProto2(SinkConfig sinkConfig)
throws IOException {
org.apache.pulsar.functions.proto.Function.FunctionDetails.Builder
functionDetailsBuilder
@@ -284,13 +262,10 @@ public void accept(String s) {
protected FunctionDetails createSinkConfig(SinkConfig sinkConfig) {
- File file = new File(jarFile);
- try {
- Reflections.loadJar(file);
- } catch (MalformedURLException e) {
- throw new RuntimeException("Failed to load user jar " + file,
e);
- }
- Class<?> typeArg = getSinkType(file);
+ // check if configs are valid
+ validateSinkConfigs(sinkConfig);
+
+ Class<?> typeArg = getSinkType(sinkConfig.getClassName());
FunctionDetails.Builder functionDetailsBuilder =
FunctionDetails.newBuilder();
if (sinkConfig.getTenant() != null) {
@@ -376,38 +351,4 @@ void runCmd() throws Exception {
print("Deleted successfully");
}
}
-
- private static SinkConfig loadSinkConfig(String file) throws IOException {
- return (SinkConfig) loadConfig(file, SinkConfig.class);
- }
-
- private static Object loadConfig(String file, Class<?> clazz) throws
IOException {
- ObjectMapper mapper = new ObjectMapper(new YAMLFactory());
- return mapper.readValue(new File(file), clazz);
- }
-
- public static boolean areAllRequiredFieldsPresentForSink(SinkConfig
sinkConfig) {
- return sinkConfig.getTenant() != null &&
!sinkConfig.getTenant().isEmpty()
- && sinkConfig.getNamespace() != null &&
!sinkConfig.getNamespace().isEmpty()
- && sinkConfig.getName() != null &&
!sinkConfig.getName().isEmpty()
- && sinkConfig.getClassName() != null &&
!sinkConfig.getClassName().isEmpty()
- && sinkConfig.getTopicToSerdeClassName() != null &&
!sinkConfig.getTopicToSerdeClassName().isEmpty()
- && sinkConfig.getParallelism() > 0;
- }
-
- private static void validateTopicName(String topic) {
- if (!TopicName.isValid(topic)) {
- throw new IllegalArgumentException(String.format("The topic name
%s is invalid", topic));
- }
- }
-
- private static ProcessingGuarantees convertProcessingGuarantee(
- FunctionConfig.ProcessingGuarantees processingGuarantees) {
- for (ProcessingGuarantees type : ProcessingGuarantees.values()) {
- if (type.name().equals(processingGuarantees.name())) {
- return type;
- }
- }
- throw new RuntimeException("Unrecognized processing guarantee: " +
processingGuarantees.name());
- }
}
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 5536bf539d..97d5fd93b4 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
@@ -19,26 +19,16 @@
package org.apache.pulsar.admin.cli;
import com.beust.jcommander.Parameter;
+import com.beust.jcommander.ParameterException;
import com.beust.jcommander.Parameters;
import com.beust.jcommander.converters.StringConverter;
-import com.fasterxml.jackson.databind.ObjectMapper;
-import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import com.google.gson.Gson;
import com.google.gson.reflect.TypeToken;
-
-import java.io.File;
-import java.io.IOException;
-import java.lang.reflect.Type;
-import java.net.MalformedURLException;
-import java.util.Map;
-
import lombok.Getter;
-
import org.apache.pulsar.client.admin.PulsarAdmin;
import org.apache.pulsar.client.admin.internal.FunctionsImpl;
import org.apache.pulsar.functions.api.utils.IdentityFunction;
import org.apache.pulsar.functions.proto.Function.FunctionDetails;
-import org.apache.pulsar.functions.proto.Function.ProcessingGuarantees;
import org.apache.pulsar.functions.proto.Function.Resources;
import org.apache.pulsar.functions.proto.Function.SinkSpec;
import org.apache.pulsar.functions.proto.Function.SourceSpec;
@@ -46,12 +36,19 @@
import org.apache.pulsar.functions.utils.Reflections;
import org.apache.pulsar.functions.utils.SourceConfig;
import org.apache.pulsar.functions.utils.Utils;
-import org.apache.pulsar.io.core.Source;
+import org.apache.pulsar.functions.utils.validation.ConfigValidation;
-import net.jodah.typetools.TypeResolver;
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Type;
+import java.net.MalformedURLException;
+import java.util.Map;
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.getSourceType;
+import static org.apache.pulsar.functions.utils.Utils.loadConfig;
@Getter
@Parameters(commandDescription = "Interface for managing Pulsar Source
(Ingress data to Pulsar)")
@@ -109,9 +106,6 @@ void runCmd() throws Exception {
public class CreateSource extends SourceCommand {
@Override
void runCmd() throws Exception {
- if (!areAllRequiredFieldsPresentForSource(sourceConfig)) {
- throw new RuntimeException("Missing arguments");
- }
admin.functions().createFunction(createSourceConfig(sourceConfig),
jarFile);
print("Created successfully");
}
@@ -121,9 +115,6 @@ void runCmd() throws Exception {
public class UpdateSource extends SourceCommand {
@Override
void runCmd() throws Exception {
- if (!areAllRequiredFieldsPresentForSource(sourceConfig)) {
- throw new RuntimeException("Missing arguments");
- }
admin.functions().updateFunction(createSourceConfig(sourceConfig),
jarFile);
print("Updated successfully");
}
@@ -145,7 +136,7 @@ void runCmd() throws Exception {
@Parameter(names = "--deserializationClassName", description = "The
classname for SerDe class for the source")
protected String deserializationClassName;
@Parameter(names = "--parallelism", description = "The source's
parallelism factor (i.e. the number of source instances to run)")
- protected String parallelism;
+ protected Integer parallelism;
@Parameter(
names = "--jar",
description = "Path to the jar file for the Source",
@@ -171,25 +162,19 @@ void processArguments() throws Exception {
super.processArguments();
if (null != sourceConfigFile) {
- this.sourceConfig = loadSourceConfig(sourceConfigFile);
+ this.sourceConfig = loadConfig(sourceConfigFile,
SourceConfig.class);
} else {
this.sourceConfig = new SourceConfig();
}
-
if (null != tenant) {
sourceConfig.setTenant(tenant);
- } else if (sourceConfig.getTenant() == null) {
- sourceConfig.setTenant(PUBLIC_TENANT);
}
if (null != namespace) {
sourceConfig.setNamespace(namespace);
- } else if (sourceConfig.getNamespace() == null) {
- sourceConfig.setNamespace(DEFAULT_NAMESPACE);
}
if (null != name) {
sourceConfig.setName(name);
}
-
if (null != className) {
this.sourceConfig.setClassName(className);
}
@@ -202,26 +187,14 @@ void processArguments() throws Exception {
if (null != processingGuarantees) {
sourceConfig.setProcessingGuarantees(processingGuarantees);
}
- if (parallelism == null) {
- if (sourceConfig.getParallelism() == 0) {
- sourceConfig.setParallelism(1);
- }
- } else {
- int num = Integer.parseInt(parallelism);
- if (num <= 0) {
- throw new IllegalArgumentException("The parallelism factor
(the number of instances) for the "
- + "connector must be positive");
- }
- sourceConfig.setParallelism(num);
+ if (parallelism != null) {
+ sourceConfig.setParallelism(parallelism);
}
if (null == jarFile) {
- throw new IllegalArgumentException("Connector JAR not
specfied");
+ throw new ParameterException("Source JAR not specfied");
}
- com.google.common.base.Preconditions.checkArgument(cpu == null ||
cpu > 0, "The cpu allocation for the source must be positive");
- com.google.common.base.Preconditions.checkArgument(ram == null ||
ram > 0, "The ram allocation for the source must be positive");
- com.google.common.base.Preconditions.checkArgument(disk == null ||
disk > 0, "The disk allocation for the source must be positive");
sourceConfig.setResources(new
org.apache.pulsar.functions.utils.Resources(cpu, ram, disk));
if (null != sourceConfigString) {
@@ -229,27 +202,36 @@ void processArguments() throws Exception {
Map<String, Object> sourceConfigMap = new
Gson().fromJson(sourceConfigString, type);
sourceConfig.setConfigs(sourceConfigMap);
}
+
+ inferMissingArguments(sourceConfig);
}
- private Class<?> getSourceType(File file) {
- if (!Reflections.classExistsInJar(file,
sourceConfig.getClassName())) {
- throw new IllegalArgumentException(String.format("Pulsar
Source class %s does not exist in jar %s",
- sourceConfig.getClassName(), jarFile));
- } else if (!Reflections.classInJarImplementsIface(file,
sourceConfig.getClassName(), Source.class)) {
- throw new IllegalArgumentException(String.format("The Pulsar
source class %s in jar %s implements does not implement " +
Source.class.getName(),
- sourceConfig.getClassName(), jarFile));
+ private void inferMissingArguments(SourceConfig sourceConfig) {
+ if (sourceConfig.getTenant() == null) {
+ sourceConfig.setTenant(PUBLIC_TENANT);
}
+ if (sourceConfig.getNamespace() == null) {
+ sourceConfig.setNamespace(DEFAULT_NAMESPACE);
+ }
+ }
- Object userClass =
Reflections.createInstance(sourceConfig.getClassName(), file);
- Class<?> typeArg;
- Source source = (Source) userClass;
- if (source == null) {
- throw new IllegalArgumentException(String.format("The Pulsar
source class %s could not be instantiated from jar %s",
- sourceConfig.getClassName(), jarFile));
+ protected void validateSourceConfigs(SourceConfig sourceConfig) {
+ File file = new File(jarFile);
+ ClassLoader userJarLoader;
+ try {
+ userJarLoader = Reflections.loadJar(file);
+ } catch (MalformedURLException e) {
+ throw new ParameterException("Failed to load user jar " + file
+ " with error " + e.getMessage());
}
- typeArg = TypeResolver.resolveRawArgument(Source.class,
source.getClass());
+ // make sure the function class loader is accessible thread-locally
+ Thread.currentThread().setContextClassLoader(userJarLoader);
- return typeArg;
+ try {
+ // Need to load jar and set context class loader before calling
+ ConfigValidation.validateConfig(sourceConfig,
FunctionConfig.Runtime.JAVA.name());
+ } catch (Exception e) {
+ throw new ParameterException(e.getMessage());
+ }
}
protected org.apache.pulsar.functions.proto.Function.FunctionDetails
createSourceConfigProto2(SourceConfig sourceConfig)
@@ -262,13 +244,10 @@ void processArguments() throws Exception {
protected FunctionDetails createSourceConfig(SourceConfig
sourceConfig) {
- File file = new File(jarFile);
- try {
- Reflections.loadJar(file);
- } catch (MalformedURLException e) {
- throw new RuntimeException("Failed to load user jar " + file,
e);
- }
- Class<?> typeArg = getSourceType(file);
+ // check if source configs are valid
+ validateSourceConfigs(sourceConfig);
+
+ Class<?> typeArg = getSourceType(sourceConfig.getClassName());
FunctionDetails.Builder functionDetailsBuilder =
FunctionDetails.newBuilder();
if (sourceConfig.getTenant() != null) {
@@ -358,33 +337,4 @@ void runCmd() throws Exception {
print("Delete source successfully");
}
}
-
- private static SourceConfig loadSourceConfig(String file) throws
IOException {
- return (SourceConfig) loadConfig(file, SourceConfig.class);
- }
-
- private static Object loadConfig(String file, Class<?> clazz) throws
IOException {
- ObjectMapper mapper = new ObjectMapper(new YAMLFactory());
- return mapper.readValue(new File(file), clazz);
- }
-
- public static boolean areAllRequiredFieldsPresentForSource(SourceConfig
sourceConfig) {
- return sourceConfig.getTenant() != null &&
!sourceConfig.getTenant().isEmpty()
- && sourceConfig.getNamespace() != null &&
!sourceConfig.getNamespace().isEmpty()
- && sourceConfig.getName() != null &&
!sourceConfig.getName().isEmpty()
- && sourceConfig.getClassName() != null &&
!sourceConfig.getClassName().isEmpty()
- && sourceConfig.getTopicName() != null &&
!sourceConfig.getTopicName().isEmpty()
- || sourceConfig.getSerdeClassName() != null
- && sourceConfig.getParallelism() > 0;
- }
-
- private static ProcessingGuarantees convertProcessingGuarantee(
- FunctionConfig.ProcessingGuarantees processingGuarantees) {
- for (ProcessingGuarantees type : ProcessingGuarantees.values()) {
- if (type.name().equals(processingGuarantees.name())) {
- return type;
- }
- }
- throw new RuntimeException("Unrecognized processing guarantee: " +
processingGuarantees.name());
- }
}
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 3f838cfdb0..9fa030771e 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
@@ -23,6 +23,14 @@
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.NotNull;
+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;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidResources;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidSinkConfig;
+import org.apache.pulsar.functions.utils.validation.ValidatorImpls;
+import org.apache.pulsar.io.core.Sink;
import java.util.HashMap;
import java.util.Map;
@@ -32,14 +40,24 @@
@Data
@EqualsAndHashCode
@ToString
+@isValidSinkConfig
public class SinkConfig {
+ @NotNull
private String tenant;
+ @NotNull
private String namespace;
+ @NotNull
private String name;
+ @NotNull
+ @isImplementationOfClass(implementsClass = Sink.class)
private String className;
+ @isMapEntryCustom(keyValidatorClasses = {
ValidatorImpls.TopicNameValidator.class },
+ valueValidatorClasses = { ValidatorImpls.SerdeValidator.class })
private Map<String, String> topicToSerdeClassName;
private Map<String, Object> configs = new HashMap<>();
+ @isPositiveNumber
private int parallelism = 1;
private FunctionConfig.ProcessingGuarantees processingGuarantees;
+ @isValidResources
private Resources resources;
}
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 89e3c802b3..cdede6434e 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
@@ -23,6 +23,14 @@
import lombok.Getter;
import lombok.Setter;
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.isImplementationOfClass;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isPositiveNumber;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidResources;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidSourceConfig;
+import
org.apache.pulsar.functions.utils.validation.ConfigValidationAnnotations.isValidTopicName;
+import org.apache.pulsar.io.core.Source;
import java.util.HashMap;
import java.util.Map;
@@ -32,15 +40,26 @@
@Data
@EqualsAndHashCode
@ToString
+@isValidSourceConfig
public class SourceConfig {
+ @NotNull
private String tenant;
+ @NotNull
private String namespace;
+ @NotNull
private String name;
+ @NotNull
+ @isImplementationOfClass(implementsClass = Source.class)
private String className;
+ @NotNull
+ @isValidTopicName
private String topicName;
+ @isImplementationOfClass(implementsClass = SerDe.class)
private String serdeClassName;
private Map<String, Object> configs = new HashMap<>();
+ @isPositiveNumber
private int parallelism = 1;
private FunctionConfig.ProcessingGuarantees processingGuarantees;
+ @isValidResources
private Resources resources;
}
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 2963e571fa..3f7c2e47dc 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
@@ -18,10 +18,23 @@
*/
package org.apache.pulsar.functions.utils;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import com.google.protobuf.AbstractMessage.Builder;
import com.google.protobuf.MessageOrBuilder;
import com.google.protobuf.util.JsonFormat;
+import lombok.AccessLevel;
+import lombok.NoArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import net.jodah.typetools.TypeResolver;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.impl.MessageIdImpl;
+import org.apache.pulsar.functions.api.Function;
+import org.apache.pulsar.functions.proto.Function.FunctionDetails.Runtime;
+import org.apache.pulsar.io.core.Sink;
+import org.apache.pulsar.io.core.Source;
+import java.io.File;
import java.io.IOException;
import java.lang.reflect.Constructor;
import java.lang.reflect.InvocationTargetException;
@@ -30,15 +43,6 @@
import java.net.ServerSocket;
import java.util.Collection;
-import lombok.AccessLevel;
-import lombok.NoArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.pulsar.client.api.MessageId;
-import org.apache.pulsar.client.impl.MessageIdImpl;
-import org.apache.pulsar.functions.api.Function;
-
-import net.jodah.typetools.TypeResolver;
-
/**
* Utils used for runtime.
*/
@@ -149,4 +153,56 @@ public static Object createInstance(String userClassName,
ClassLoader classLoade
return result;
}
+
+ public static <T> T loadConfig(String file, Class<T> clazz) throws
IOException {
+ ObjectMapper mapper = new ObjectMapper(new YAMLFactory());
+ return mapper.readValue(new File(file), clazz);
+ }
+
+ public static Runtime convertRuntime(FunctionConfig.Runtime runtime) {
+ for (Runtime type : Runtime.values()) {
+ if (type.name().equals(runtime.name())) {
+ return type;
+ }
+ }
+ throw new RuntimeException("Unrecognized runtime: " + runtime.name());
+ }
+
+ public static
org.apache.pulsar.functions.proto.Function.ProcessingGuarantees
convertProcessingGuarantee(
+ FunctionConfig.ProcessingGuarantees processingGuarantees) {
+ for (org.apache.pulsar.functions.proto.Function.ProcessingGuarantees
type :
org.apache.pulsar.functions.proto.Function.ProcessingGuarantees.values()) {
+ if (type.name().equals(processingGuarantees.name())) {
+ return type;
+ }
+ }
+ throw new RuntimeException("Unrecognized processing guarantee: " +
processingGuarantees.name());
+ }
+
+ public static Class<?> getSourceType(String className) {
+
+ Object userClass = Reflections.createInstance(className,
Thread.currentThread().getContextClassLoader());
+ Class<?> typeArg;
+ Source source = (Source) userClass;
+ if (source == null) {
+ throw new IllegalArgumentException(String.format("The Pulsar
source class %s could not be instantiated",
+ className));
+ }
+ typeArg = TypeResolver.resolveRawArgument(Source.class,
source.getClass());
+
+ return typeArg;
+ }
+
+ public static Class<?> getSinkType(String className) {
+
+ Object userClass = Reflections.createInstance(className,
Thread.currentThread().getContextClassLoader());
+ Class<?> typeArg;
+ Sink sink = (Sink) userClass;
+ if (sink == null) {
+ throw new IllegalArgumentException(String.format("The Pulsar sink
class %s could not be instantiated",
+ className));
+ }
+ typeArg = TypeResolver.resolveRawArgument(Sink.class, sink.getClass());
+
+ return typeArg;
+ }
}
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 b01331e53e..934f6f55e2 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
@@ -197,6 +197,17 @@
Class<?> validatorClass() default
ValidatorImpls.FunctionConfigValidator.class;
}
+ @Retention(RetentionPolicy.RUNTIME)
+ @Target({ElementType.TYPE})
+ public @interface isValidSourceConfig {
+ Class<?> validatorClass() default
ValidatorImpls.SourceConfigValidator.class;
+ }
+
+ @Retention(RetentionPolicy.RUNTIME)
+ @Target({ElementType.TYPE})
+ public @interface isValidSinkConfig {
+ Class<?> validatorClass() default
ValidatorImpls.SinkConfigValidator.class;
+ }
/**
* Field names for annotations
*/
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 bc86c2d42a..38b54abaa9 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
@@ -21,12 +21,15 @@
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import net.jodah.typetools.TypeResolver;
+import org.apache.commons.lang.StringUtils;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.functions.api.SerDe;
import org.apache.pulsar.functions.api.utils.DefaultSerDe;
import org.apache.pulsar.functions.utils.FunctionConfig;
import org.apache.pulsar.functions.utils.Reflections;
import org.apache.pulsar.functions.utils.Resources;
+import org.apache.pulsar.functions.utils.SinkConfig;
+import org.apache.pulsar.functions.utils.SourceConfig;
import org.apache.pulsar.functions.utils.Utils;
import org.apache.pulsar.functions.utils.WindowConfig;
@@ -36,6 +39,9 @@
import java.util.HashSet;
import java.util.Map;
+import static org.apache.pulsar.functions.utils.Utils.getSinkType;
+import static org.apache.pulsar.functions.utils.Utils.getSourceType;
+
@Slf4j
public class ValidatorImpls {
/**
@@ -181,16 +187,21 @@ public void validateField(String name, Object o) {
}
SimpleTypeValidator.validateField(name, String.class, o);
String className = (String) o;
+ if (StringUtils.isEmpty(className)) {
+ return;
+ }
+
+ Class<?> objectClass;
try {
- ClassLoader clsLoader =
Thread.currentThread().getContextClassLoader();
- Class<?> objectClass = clsLoader.loadClass(className);
- if (!this.classImplements.isAssignableFrom(objectClass)) {
- throw new IllegalArgumentException(
- String.format("Field '%s' with value '%s' does not
implement %s ",
- name, o, this.classImplements.getName()));
- }
+ objectClass = loadClass(className);
} catch (ClassNotFoundException e) {
- throw new RuntimeException(e);
+ throw new IllegalArgumentException("Cannot find/load class " +
className);
+ }
+
+ if (!this.classImplements.isAssignableFrom(objectClass)) {
+ throw new IllegalArgumentException(
+ String.format("Field '%s' with value '%s' does not
implement %s ",
+ name, o, this.classImplements.getName()));
}
}
}
@@ -213,6 +224,9 @@ public void validateField(String name, Object o) {
}
SimpleTypeValidator.validateField(name, String.class, o);
String className = (String) o;
+ if (StringUtils.isEmpty(className)) {
+ return;
+ }
int count = 0;
for (Class<?> classImplements : classesImplements) {
Class<?> objectClass = null;
@@ -229,7 +243,7 @@ public void validateField(String name, Object o) {
if (count == 0) {
throw new IllegalArgumentException(
String.format("Field '%s' with value '%s' does not
implement any of these classes %s",
- name, o, classesImplements));
+ name, o, Arrays.asList(classesImplements)));
}
}
}
@@ -337,7 +351,6 @@ private static void doJavaChecks(FunctionConfig
functionConfig, String name) {
// implements SerDe class
functionConfig.getCustomSerdeInputs().forEach((topicName,
inputSerializer) -> {
-
Class<?> serdeClass;
try {
serdeClass = loadClass(inputSerializer);
@@ -603,7 +616,133 @@ public void validateField(String name, Object o) {
}
}
- /**
+ public static class SourceConfigValidator extends Validator {
+ @Override
+ public void validateField(String name, Object o) {
+ SourceConfig sourceConfig = (SourceConfig) o;
+ Class<?> typeArg = getSourceType(sourceConfig.getClassName());
+ String serdeClassname = sourceConfig.getSerdeClassName();
+
+ ClassLoader clsLoader =
Thread.currentThread().getContextClassLoader();
+
+ if (StringUtils.isEmpty(serdeClassname)) {
+ serdeClassname = DefaultSerDe.class.getName();
+ }
+
+ try {
+ loadClass(serdeClassname);
+ } catch (ClassNotFoundException e) {
+ throw new IllegalArgumentException(
+ String.format("The input serialization/deserialization
class %s does not exist", serdeClassname));
+
+ }
+
+ try {
+ new
ValidatorImpls.ImplementsClassValidator(SerDe.class).validateField(name,
serdeClassname);
+ } catch (IllegalArgumentException ex) {
+ throw new IllegalArgumentException(
+ String.format("The input serialization/deserialization
class %s does not not implement %s",
+ serdeClassname,
SerDe.class.getCanonicalName()));
+ }
+
+ if (serdeClassname.equals(DefaultSerDe.class.getName())) {
+ if (!DefaultSerDe.IsSupportedType(typeArg)) {
+ throw new IllegalArgumentException("The default Serializer
does not support type " + typeArg);
+ }
+ } else {
+ SerDe serDe = (SerDe)
Reflections.createInstance(serdeClassname, clsLoader);
+ if (serDe == null) {
+ throw new IllegalArgumentException(String.format("The
SerDe class %s does not exist",
+ serdeClassname));
+ }
+ Class<?>[] serDeTypes =
TypeResolver.resolveRawArguments(SerDe.class, serDe.getClass());
+
+ // type inheritance information seems to be lost in generic
type
+ // load the actual type class for verification
+ Class<?> fnInputClass;
+ Class<?> serdeInputClass;
+ try {
+ fnInputClass = Class.forName(typeArg.getName(), true,
clsLoader);
+ // get output serde
+ serdeInputClass = Class.forName(serDeTypes[1].getName(),
true, clsLoader);
+ } catch (ClassNotFoundException e) {
+ throw new IllegalArgumentException("Failed to load type
class", e);
+ }
+
+ if (!fnInputClass.isAssignableFrom(serdeInputClass)) {
+ throw new IllegalArgumentException("Serializer type
mismatch " + typeArg + " vs " +
+ serDeTypes[1]);
+ }
+ }
+ }
+ }
+
+ public static class SinkConfigValidator extends Validator {
+ @Override
+ public void validateField(String name, Object o) {
+ SinkConfig sinkConfig = (SinkConfig) o;
+ Class<?> typeArg = getSinkType(sinkConfig.getClassName());
+
+ ClassLoader clsLoader =
Thread.currentThread().getContextClassLoader();
+
+ sinkConfig.getTopicToSerdeClassName().forEach((topicName,
serdeClassname) -> {
+ if (StringUtils.isEmpty(serdeClassname)) {
+ serdeClassname = DefaultSerDe.class.getName();
+ }
+
+ try {
+ loadClass(serdeClassname);
+ } catch (ClassNotFoundException e) {
+ throw new IllegalArgumentException(
+ String.format("The input
serialization/deserialization class %s does not exist", serdeClassname));
+ }
+
+ try {
+ new
ValidatorImpls.ImplementsClassValidator(SerDe.class).validateField(name,
serdeClassname);
+ } catch (IllegalArgumentException ex) {
+ throw new IllegalArgumentException(
+ String.format("The input
serialization/deserialization class %s does not not " +
+ "implement %s",
+ serdeClassname,
SerDe.class.getCanonicalName()));
+ }
+
+ if (serdeClassname.equals(DefaultSerDe.class.getName())) {
+ if (!DefaultSerDe.IsSupportedType(typeArg)) {
+ throw new IllegalArgumentException("The default
Serializer does not support type " +
+ typeArg);
+ }
+ } else {
+ SerDe serDe = (SerDe)
Reflections.createInstance(serdeClassname, clsLoader);
+ if (serDe == null) {
+ throw new IllegalArgumentException(String.format("The
SerDe class %s does not exist",
+ serdeClassname));
+ }
+ Class<?>[] serDeTypes =
TypeResolver.resolveRawArguments(SerDe.class, serDe.getClass());
+
+ // type inheritance information seems to be lost in
generic type
+ // load the actual type class for verification
+ Class<?> fnInputClass;
+ Class<?> serdeInputClass;
+ try {
+ fnInputClass = Class.forName(typeArg.getName(), true,
clsLoader);
+ // get input serde
+ serdeInputClass =
Class.forName(serDeTypes[0].getName(), true, clsLoader);
+ } catch (ClassNotFoundException e) {
+ throw new IllegalArgumentException("Failed to load
type class", e);
+ }
+
+ if (!fnInputClass.isAssignableFrom(serdeInputClass)) {
+ throw new IllegalArgumentException("Serializer type
mismatch " + typeArg + " vs " +
+ serDeTypes[0]);
+ }
+ }
+ });
+ }
+ }
+
+
+
+ /**
* 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