This is an automated email from the ASF dual-hosted git repository.

mmerli 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 d6a341f  refactor source and sink classname to be pulsar source and 
sink when not set (#1769)
d6a341f is described below

commit d6a341f4b5160c30c55e071b1d2e6d1860649c07
Author: Boyang Jerry Peng <[email protected]>
AuthorDate: Mon May 14 14:05:25 2018 -0700

    refactor source and sink classname to be pulsar source and sink when not 
set (#1769)
---
 .../main/java/org/apache/pulsar/admin/cli/CmdFunctions.java  |  6 ------
 .../src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java  |  2 +-
 .../main/java/org/apache/pulsar/admin/cli/CmdSources.java    |  4 ++--
 .../pulsar/functions/instance/JavaInstanceRunnable.java      | 11 ++++++-----
 .../apache/pulsar/functions/runtime/JavaInstanceMain.java    | 12 ++++++++----
 5 files changed, 17 insertions(+), 18 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 9f2ce10..b704ec7 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
@@ -651,9 +651,6 @@ public class CmdFunctions extends CmdBase {
             Map<String, String> topicToSerDeClassNameMap = new HashMap<>();
             
topicToSerDeClassNameMap.putAll(functionConfig.getCustomSerdeInputs());
             SourceSpec.Builder sourceSpecBuilder = SourceSpec.newBuilder();
-            if (functionConfig.getRuntime() == FunctionConfig.Runtime.JAVA) {
-                sourceSpecBuilder.setClassName(PulsarSource.class.getName());
-            }
             functionConfig.getInputs().forEach(v -> 
topicToSerDeClassNameMap.put(v, ""));
             
sourceSpecBuilder.putAllTopicsToSerDeClassName(topicToSerDeClassNameMap);
             if (functionConfig.getSubscriptionType() != null) {
@@ -667,9 +664,6 @@ public class CmdFunctions extends CmdBase {
 
             // Setup sink
             SinkSpec.Builder sinkSpecBuilder = SinkSpec.newBuilder();
-            if (functionConfig.getRuntime() == FunctionConfig.Runtime.JAVA) {
-                sinkSpecBuilder.setClassName(PulsarSink.class.getName());
-            }
             if (functionConfig.getOutput() != null) {
                 sinkSpecBuilder.setTopic(functionConfig.getOutput());
             }
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 dfc7d0b..fb0bd65 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
@@ -288,8 +288,8 @@ public class CmdSinks extends CmdBase {
             }
 
             // set source spec
+            // source spec classname should be empty so that the default 
pulsar source will be used
             SourceSpec.Builder sourceSpecBuilder = SourceSpec.newBuilder();
-            sourceSpecBuilder.setClassName(PulsarSource.class.getName());
             
sourceSpecBuilder.setSubscriptionType(Function.SubscriptionType.SHARED);
             
sourceSpecBuilder.putAllTopicsToSerDeClassName(sinkConfig.getTopicToSerdeClassName());
             sourceSpecBuilder.setTypeClassName(typeArg.getName());
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 267673c..b78ebde 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
@@ -279,9 +279,9 @@ public class CmdSources extends CmdBase {
             sourceSpecBuilder.setTypeClassName(typeArg.getName());
             functionDetailsBuilder.setSource(sourceSpecBuilder);
 
-            // set up sink spec
+            // set up sink spec.
+            // Sink spec classname should be empty so that the default pulsar 
sink will be used
             SinkSpec.Builder sinkSpecBuilder = SinkSpec.newBuilder();
-            sinkSpecBuilder.setClassName(PulsarSink.class.getName());
             if (sourceConfig.getSerdeClassName() != null && 
!sourceConfig.getSerdeClassName().isEmpty()) {
                 
sinkSpecBuilder.setSerDeClassName(sourceConfig.getSerdeClassName());
             }
diff --git 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
index 4257053..e712dab 100644
--- 
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
+++ 
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/JavaInstanceRunnable.java
@@ -455,8 +455,8 @@ public class JavaInstanceRunnable implements AutoCloseable, 
Runnable {
 
         SourceSpec sourceSpec = 
this.instanceConfig.getFunctionDetails().getSource();
         Object object;
-        if (sourceSpec.getClassName().equals(PulsarSource.class.getName())) {
-
+        // If source classname is not set, we default pulsar source
+        if (sourceSpec.getClassName().isEmpty()) {
             PulsarSourceConfig pulsarSourceConfig = new PulsarSourceConfig();
             
pulsarSourceConfig.setTopicSerdeClassNameMap(sourceSpec.getTopicsToSerDeClassNameMap());
             pulsarSourceConfig.setSubscriptionName(
@@ -472,7 +472,7 @@ public class JavaInstanceRunnable implements AutoCloseable, 
Runnable {
             Class[] paramTypes = {PulsarClient.class, 
PulsarSourceConfig.class};
 
             object = Reflections.createInstance(
-                    sourceSpec.getClassName(),
+                    PulsarSource.class.getName(),
                     PulsarSource.class.getClassLoader(), params, paramTypes);
 
         } else {
@@ -498,7 +498,8 @@ public class JavaInstanceRunnable implements AutoCloseable, 
Runnable {
 
         SinkSpec sinkSpec = this.instanceConfig.getFunctionDetails().getSink();
         Object object;
-        if (sinkSpec.getClassName().equals(PulsarSink.class.getName())) {
+        // If sink classname is not set, we default pulsar sink
+        if (sinkSpec.getClassName().isEmpty()) {
             PulsarSinkConfig pulsarSinkConfig = new PulsarSinkConfig();
             
pulsarSinkConfig.setProcessingGuarantees(FunctionConfig.ProcessingGuarantees.valueOf(
                     
this.instanceConfig.getFunctionDetails().getProcessingGuarantees().name()));
@@ -510,7 +511,7 @@ public class JavaInstanceRunnable implements AutoCloseable, 
Runnable {
             Class[] paramTypes = {PulsarClient.class, PulsarSinkConfig.class};
 
             object = Reflections.createInstance(
-                    sinkSpec.getClassName(),
+                    PulsarSink.class.getName(),
                     PulsarSink.class.getClassLoader(), params, paramTypes);
         } else {
             object = Reflections.createInstance(
diff --git 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceMain.java
 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceMain.java
index 1f34b8b..4e23038 100644
--- 
a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceMain.java
+++ 
b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceMain.java
@@ -92,7 +92,7 @@ public class JavaInstanceMain {
     @Parameter(names = "--auto_ack", description = "Enable Auto Acking?\n")
     protected String autoAck = "true";
 
-    @Parameter(names = "--source_classname", description = "The source 
classname", required = true)
+    @Parameter(names = "--source_classname", description = "The source 
classname")
     protected String sourceClassname;
 
     @Parameter(names = "--source_configs", description = "The source configs")
@@ -113,7 +113,7 @@ public class JavaInstanceMain {
     @Parameter(names = "--sink_configs", description = "The sink configs\n")
     protected String sinkConfigs;
 
-    @Parameter(names = "--sink_classname", description = "The sink 
classname\n", required = true)
+    @Parameter(names = "--sink_classname", description = "The sink 
classname\n")
     protected String sinkClassname;
 
     @Parameter(names = "--sink_topic", description = "The sink Topic Name\n")
@@ -154,7 +154,9 @@ public class JavaInstanceMain {
 
         // Setup source
         SourceSpec.Builder sourceDetailsBuilder = SourceSpec.newBuilder();
-        sourceDetailsBuilder.setClassName(sourceClassname);
+        if (sourceClassname != null) {
+            sourceDetailsBuilder.setClassName(sourceClassname);
+        }
         if (sourceConfigs != null && !sourceConfigs.isEmpty()) {;
             sourceDetailsBuilder.setConfigs(sourceConfigs);
         }
@@ -165,7 +167,9 @@ public class JavaInstanceMain {
 
         // Setup sink
         SinkSpec.Builder sinkSpecBuilder = SinkSpec.newBuilder();
-        sinkSpecBuilder.setClassName(sinkClassname);
+        if (sinkClassname != null) {
+            sinkSpecBuilder.setClassName(sinkClassname);
+        }
         if (sinkConfigs != null) {
             sinkSpecBuilder.setConfigs(sinkConfigs);
         }

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to