This is an automated email from the ASF dual-hosted git repository.
sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 6be584f improving source and sink validation (#1897)
6be584f is described below
commit 6be584f6a35e727377f0bc0968de57c50478164d
Author: Boyang Jerry Peng <[email protected]>
AuthorDate: Sun Jun 3 14:37:00 2018 -0700
improving source and sink validation (#1897)
* add source config validation
* improving source and sink validation
* cleaning up code
* refactor sink infer config code
---
.../org/apache/pulsar/admin/cli/CmdFunctions.java | 42 +-----
.../java/org/apache/pulsar/admin/cli/CmdSinks.java | 157 +++++++-------------
.../org/apache/pulsar/admin/cli/CmdSources.java | 136 ++++++------------
.../apache/pulsar/functions/utils/SinkConfig.java | 18 +++
.../pulsar/functions/utils/SourceConfig.java | 19 +++
.../org/apache/pulsar/functions/utils/Utils.java | 74 ++++++++--
.../validation/ConfigValidationAnnotations.java | 11 ++
.../functions/utils/validation/ValidatorImpls.java | 160 +++++++++++++++++++--
8 files changed, 361 insertions(+), 256 deletions(-)
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 8590dd5..30916cd 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.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.common.annotations.VisibleForTesting;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
@@ -46,7 +44,6 @@ import org.apache.pulsar.client.admin.internal.FunctionsImpl;
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;
@@ -261,7 +258,7 @@ public class CmdFunctions extends CmdBase {
// 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();
}
@@ -477,6 +474,9 @@ public class CmdFunctions extends CmdBase {
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
@@ -544,11 +544,11 @@ public class CmdFunctions extends CmdBase {
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<>();
@@ -609,8 +609,6 @@ public class CmdFunctions extends CmdBase {
@Override
void runCmd() throws Exception {
- // check if function configs are valid
- validateFunctionConfigs(functionConfig);
CmdFunctions.startLocalRun(convertProto2(functionConfig),
functionConfig.getParallelism(), brokerServiceUrl,
userCodeFile, admin);
}
@@ -620,8 +618,6 @@ public class CmdFunctions extends CmdBase {
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");
}
@@ -660,8 +656,6 @@ public class CmdFunctions extends CmdBase {
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");
}
@@ -869,30 +863,6 @@ public class CmdFunctions extends CmdBase {
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 2579c41..fdbe291 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.FunctionConfig;
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 @@ public class CmdSinks extends CmdBase {
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 @@ public class CmdSinks extends CmdBase {
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 @@ public class CmdSinks extends CmdBase {
@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 @@ public class CmdSinks extends CmdBase {
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 @@ public class CmdSinks extends CmdBase {
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 class CmdSinks extends CmdBase {
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 class CmdSinks extends CmdBase {
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 @@ public class CmdSinks extends CmdBase {
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 5536bf5..97d5fd9 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.FunctionConfig;
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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
@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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
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 @@ public class CmdSources extends CmdBase {
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 3f838cf..9fa0307 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.EqualsAndHashCode;
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 @@ import java.util.Map;
@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 89e3c80..cdede64 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.EqualsAndHashCode;
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 @@ import java.util.Map;
@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 2963e57..3f7c2e4 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.lang.reflect.Type;
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 class Utils {
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 b01331e..934f6f5 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 @@ public class ConfigValidationAnnotations {
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 5bebf19..3e0a168 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 @@ package org.apache.pulsar.functions.utils.validation;
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.Collection;
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 class ValidatorImpls {
}
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 class ValidatorImpls {
}
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 class ValidatorImpls {
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)));
}
}
}
@@ -609,7 +623,133 @@ public class ValidatorImpls {
}
}
- /**
+ 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 {
--
To stop receiving notification emails like this one, please contact
[email protected].