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].

Reply via email to