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