merlimat closed pull request #1781: fixing behavior when configs are empty
URL: https://github.com/apache/incubator-pulsar/pull/1781
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
index b704ec7361..cf28b5a269 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
@@ -306,6 +306,15 @@ void processArguments() throws Exception {
Map<String, Object> userConfigMap = new
Gson().fromJson(userConfigString, type);
functionConfig.setUserConfig(userConfigMap);
}
+ if (functionConfig.getInputs() == null) {
+ functionConfig.setInputs(new LinkedList<>());
+ }
+ if (functionConfig.getCustomSerdeInputs() == null) {
+ functionConfig.setCustomSerdeInputs(new HashMap<>());
+ }
+ if (functionConfig.getUserConfig() == null) {
+ functionConfig.setUserConfig(new HashMap<>());
+ }
if (functionConfig.getInputs().isEmpty() &&
functionConfig.getCustomSerdeInputs().isEmpty()) {
throw new RuntimeException("No input topic(s) specified for
the function");
@@ -648,11 +657,12 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
FunctionDetails.Builder functionDetailsBuilder =
FunctionDetails.newBuilder();
// Setup source
+ SourceSpec.Builder sourceSpecBuilder = SourceSpec.newBuilder();
Map<String, String> topicToSerDeClassNameMap = new HashMap<>();
topicToSerDeClassNameMap.putAll(functionConfig.getCustomSerdeInputs());
- SourceSpec.Builder sourceSpecBuilder = SourceSpec.newBuilder();
functionConfig.getInputs().forEach(v ->
topicToSerDeClassNameMap.put(v, ""));
sourceSpecBuilder.putAllTopicsToSerDeClassName(topicToSerDeClassNameMap);
+
if (functionConfig.getSubscriptionType() != null) {
sourceSpecBuilder
.setSubscriptionType(convertSubscriptionType(functionConfig.getSubscriptionType()));
@@ -697,6 +707,7 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
Map<String, Object> configs = new HashMap<>();
configs.putAll(functionConfig.getUserConfig());
+
// windowing related
WindowConfig windowConfig = functionConfig.getWindowConfig();
if (windowConfig != null) {
@@ -710,7 +721,9 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
functionDetailsBuilder.setClassName(functionConfig.getClassName());
}
}
- functionDetailsBuilder.setUserConfig(new Gson().toJson(configs));
+ if (!configs.isEmpty()) {
+ functionDetailsBuilder.setUserConfig(new
Gson().toJson(configs));
+ }
functionDetailsBuilder.setAutoAck(functionConfig.isAutoAck());
functionDetailsBuilder.setParallelism(functionConfig.getParallelism());
diff --git
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
index d1971d0879..9eddc690e4 100644
---
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
+++
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/ContextImpl.java
@@ -113,8 +113,12 @@ public ContextImpl(InstanceConfig config, Logger logger,
PulsarClient client,
producerConfiguration.setBatchingEnabled(true);
producerConfiguration.setBatchingMaxPublishDelay(1,
TimeUnit.MILLISECONDS);
producerConfiguration.setMaxPendingMessages(1000000);
- userConfigs = new
Gson().fromJson(config.getFunctionDetails().getUserConfig(),
- new TypeToken<Map<String, Object>>(){}.getType());
+ if (config.getFunctionDetails().getUserConfig().isEmpty()) {
+ userConfigs = new HashMap<>();
+ } else {
+ userConfigs = new
Gson().fromJson(config.getFunctionDetails().getUserConfig(),
+ new TypeToken<Map<String, Object>>(){}.getType());
+ }
}
public void setCurrentMessageContext(MessageId messageId, String
topicName) {
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 e712dabaa8..4aaed5bab3 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
@@ -28,6 +28,7 @@
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
@@ -490,8 +491,12 @@ public void setupInput() throws Exception {
}
this.source = (Source) object;
- this.source.open(new Gson().fromJson(sourceSpec.getConfigs(),
- new TypeToken<Map<String, Object>>(){}.getType()));
+ if (sourceSpec.getConfigs().isEmpty()) {
+ this.source.open(new HashMap<>());
+ } else {
+ this.source.open(new Gson().fromJson(sourceSpec.getConfigs(),
+ new TypeToken<Map<String, Object>>(){}.getType()));
+ }
}
public void setupOutput() throws Exception {
@@ -526,6 +531,11 @@ public void setupOutput() throws Exception {
} else {
throw new RuntimeException("Sink does not implement correct
interface");
}
- this.sink.open(new Gson().fromJson(sinkSpec.getConfigs(), new
TypeToken<Map<String, Object>>(){}.getType()));
+ if (sinkSpec.getConfigs().isEmpty()) {
+ this.sink.open(new HashMap<>());
+ } else {
+ this.sink.open(new Gson().fromJson(sinkSpec.getConfigs(),
+ new TypeToken<Map<String, Object>>() {}.getType()));
+ }
}
}
diff --git a/pulsar-functions/instance/src/main/python/contextimpl.py
b/pulsar-functions/instance/src/main/python/contextimpl.py
index e17b296a49..3463d7ba4d 100644
--- a/pulsar-functions/instance/src/main/python/contextimpl.py
+++ b/pulsar-functions/instance/src/main/python/contextimpl.py
@@ -60,7 +60,9 @@ def __init__(self, instance_config, logger, pulsar_client,
user_code, consumers)
self.current_message_id = None
self.current_input_topic_name = None
self.current_start_time = None
- self.user_config = json.loads(instance_config.function_details.userConfig);
+ self.user_config = json.loads(instance_config.function_details.userConfig)
\
+ if instance_config.function_details.userConfig \
+ else []
# Called on a per message basis to set the context for the current message
def set_current_message_context(self, msgid, topic):
diff --git
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
index 5895179a8f..ffdda63e3f 100644
---
a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
+++
b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfig.java
@@ -66,15 +66,15 @@
private String name;
private String className;
- private Collection<String> inputs = new LinkedList<>();
- private Map<String, String> customSerdeInputs = new HashMap<>();
+ private Collection<String> inputs;
+ private Map<String, String> customSerdeInputs;
private String output;
private String outputSerdeClassName;
private String logTopic;
private ProcessingGuarantees processingGuarantees;
- private Map<String, Object> userConfig = new HashMap<>();
+ private Map<String, Object> userConfig;
private SubscriptionType subscriptionType;
private Runtime runtime;
private boolean autoAck;
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services