jerrypeng closed pull request #1880: removing allowing users set subscription
type
URL: https://github.com/apache/incubator-pulsar/pull/1880
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 c4648e9d88..b3e4f6bb21 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
@@ -230,8 +230,6 @@ void processArguments() throws Exception {
protected String fnConfigFile;
@Parameter(names = "--processingGuarantees", description = "The
processing guarantees (aka delivery semantics) applied to the function")
protected FunctionConfig.ProcessingGuarantees processingGuarantees;
- @Parameter(names = "--subscriptionType", description = "The type of
subscription used by the function when consuming messages from the input
topic(s)")
- protected FunctionConfig.SubscriptionType subscriptionType;
@Parameter(names = "--userConfig", description = "User-defined config
key/values")
protected String userConfigString;
@Parameter(names = "--parallelism", description = "The function's
parallelism factor (i.e. the number of function instances to run)")
@@ -304,9 +302,6 @@ void processArguments() throws Exception {
if (null != processingGuarantees) {
functionConfig.setProcessingGuarantees(processingGuarantees);
}
- if (null != subscriptionType) {
- functionConfig.setSubscriptionType(subscriptionType);
- }
if (null != userConfigString) {
Type type = new TypeToken<Map<String, String>>(){}.getType();
Map<String, Object> userConfigMap = new
Gson().fromJson(userConfigString, type);
@@ -491,10 +486,24 @@ protected FunctionDetails convert(FunctionConfig
functionConfig)
functionConfig.getInputs().forEach(v ->
topicToSerDeClassNameMap.put(v, ""));
sourceSpecBuilder.putAllTopicsToSerDeClassName(topicToSerDeClassNameMap);
- if (functionConfig.getSubscriptionType() != null) {
- sourceSpecBuilder
-
.setSubscriptionType(convertSubscriptionType(functionConfig.getSubscriptionType()));
+ // Set subscription type based on processing semantics
+ if (functionConfig.getProcessingGuarantees() != null) {
+ switch (functionConfig.getProcessingGuarantees()) {
+ case ATMOST_ONCE:
+
sourceSpecBuilder.setSubscriptionType(SubscriptionType.SHARED);
+ break;
+ case ATLEAST_ONCE:
+
sourceSpecBuilder.setSubscriptionType(SubscriptionType.SHARED);
+ break;
+ case EFFECTIVELY_ONCE:
+
sourceSpecBuilder.setSubscriptionType(SubscriptionType.FAILOVER);
+ break;
+ default:
+ throw new RuntimeException("Unknown processing
guarantee: "
+ +
functionConfig.getProcessingGuarantees().name());
+ }
}
+
if (typeArgs != null) {
sourceSpecBuilder.setTypeClassName(typeArgs[0].getName());
}
@@ -865,16 +874,6 @@ private static FunctionConfig loadConfig(File file) throws
IOException {
throw new RuntimeException("Unrecognized runtime: " + runtime.name());
}
- private static SubscriptionType convertSubscriptionType(
- FunctionConfig.SubscriptionType subscriptionType) {
- for (SubscriptionType type : SubscriptionType.values()) {
- if (type.name().equals(subscriptionType.name())) {
- return type;
- }
- }
- throw new RuntimeException("Unrecognized subscription type: " +
subscriptionType.name());
- }
-
private static ProcessingGuarantees convertProcessingGuarantee(
FunctionConfig.ProcessingGuarantees processingGuarantees) {
for (ProcessingGuarantees type : ProcessingGuarantees.values()) {
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 28e91f3c3b..10c93ec50e 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
@@ -51,6 +51,7 @@
import org.apache.logging.log4j.core.config.LoggerConfig;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.SubscriptionType;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.functions.api.Function;
import org.apache.pulsar.functions.proto.InstanceCommunication;
@@ -468,8 +469,16 @@ public void setupInput() throws Exception {
pulsarSourceConfig.setProcessingGuarantees(
FunctionConfig.ProcessingGuarantees.valueOf(
this.instanceConfig.getFunctionDetails().getProcessingGuarantees().name()));
- pulsarSourceConfig.setSubscriptionType(
-
FunctionConfig.SubscriptionType.valueOf(sourceSpec.getSubscriptionType().name()));
+
+ switch (sourceSpec.getSubscriptionType()) {
+ case FAILOVER:
+
pulsarSourceConfig.setSubscriptionType(SubscriptionType.Failover);
+ break;
+ default:
+
pulsarSourceConfig.setSubscriptionType(SubscriptionType.Shared);
+ break;
+ }
+
pulsarSourceConfig.setTypeClassName(sourceSpec.getTypeClassName());
Object[] params = {this.client, pulsarSourceConfig};
diff --git
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/sink/PulsarSinkConfig.java
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/sink/PulsarSinkConfig.java
index 60baa1a5b1..6f0385aeeb 100644
---
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/sink/PulsarSinkConfig.java
+++
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/sink/PulsarSinkConfig.java
@@ -28,7 +28,6 @@
@ToString
public class PulsarSinkConfig {
private FunctionConfig.ProcessingGuarantees processingGuarantees;
- private FunctionConfig.SubscriptionType subscriptionType;
private String topic;
private String serDeClassName;
private String typeClassName;
diff --git
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSource.java
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSource.java
index 25b0874446..c27bda8378 100644
---
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSource.java
+++
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSource.java
@@ -63,7 +63,7 @@ public void open(Map<String, Object> config) throws Exception
{
this.inputConsumer = this.pulsarClient.newConsumer()
.topics(new
ArrayList<>(this.pulsarSourceConfig.getTopicSerdeClassNameMap().keySet()))
.subscriptionName(this.pulsarSourceConfig.getSubscriptionName())
-
.subscriptionType(this.pulsarSourceConfig.getSubscriptionType().get())
+
.subscriptionType(this.pulsarSourceConfig.getSubscriptionType())
.ackTimeout(1, TimeUnit.MINUTES)
.subscribe();
}
diff --git
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSourceConfig.java
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSourceConfig.java
index 4d5e540e1e..95c1001fe2 100644
---
a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSourceConfig.java
+++
b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarSourceConfig.java
@@ -22,6 +22,8 @@
import lombok.Getter;
import lombok.Setter;
import lombok.ToString;
+import org.apache.pulsar.client.api.SubscriptionType;
+import org.apache.pulsar.functions.proto.Function;
import org.apache.pulsar.functions.utils.FunctionConfig;
import java.io.IOException;
@@ -33,7 +35,7 @@
public class PulsarSourceConfig {
private FunctionConfig.ProcessingGuarantees processingGuarantees;
- private FunctionConfig.SubscriptionType subscriptionType;
+ SubscriptionType subscriptionType;
private String subscriptionName;
private Map<String, String> topicSerdeClassNameMap;
private String typeClassName;
diff --git
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/sink/PulsarSinkTest.java
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/sink/PulsarSinkTest.java
index 86ea562b5e..df4e83b28b 100644
---
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/sink/PulsarSinkTest.java
+++
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/sink/PulsarSinkTest.java
@@ -81,7 +81,6 @@ private static PulsarClient getPulsarClient() throws
PulsarClientException {
private static PulsarSinkConfig getPulsarConfigs() {
PulsarSinkConfig pulsarConfig = new PulsarSinkConfig();
pulsarConfig.setProcessingGuarantees(FunctionConfig.ProcessingGuarantees.ATLEAST_ONCE);
-
pulsarConfig.setSubscriptionType(FunctionConfig.SubscriptionType.FAILOVER);
pulsarConfig.setTopic(TOPIC);
pulsarConfig.setSerDeClassName(serDeClassName);
pulsarConfig.setTypeClassName(String.class.getName());
diff --git
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/source/PulsarSourceTest.java
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/source/PulsarSourceTest.java
index 77d397cacd..4c81016aa7 100644
---
a/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/source/PulsarSourceTest.java
+++
b/pulsar-functions/instance/src/test/java/org/apache/pulsar/functions/source/PulsarSourceTest.java
@@ -87,7 +87,6 @@ private static PulsarClient getPulsarClient() throws
PulsarClientException {
private static PulsarSourceConfig getPulsarConfigs() {
PulsarSourceConfig pulsarConfig = new PulsarSourceConfig();
pulsarConfig.setProcessingGuarantees(FunctionConfig.ProcessingGuarantees.ATLEAST_ONCE);
-
pulsarConfig.setSubscriptionType(FunctionConfig.SubscriptionType.FAILOVER);
pulsarConfig.setTopicSerdeClassNameMap(topicSerdeClassNameMap);
pulsarConfig.setTypeClassName(String.class.getName());
return pulsarConfig;
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 2c97e65bbb..742f5309a2 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
@@ -55,20 +55,6 @@
EFFECTIVELY_ONCE
}
- public enum SubscriptionType {
- SHARED,
- FAILOVER;
-
- public org.apache.pulsar.client.api.SubscriptionType get() {
- switch (this) {
- case FAILOVER:
- return
org.apache.pulsar.client.api.SubscriptionType.Failover;
- default:
- return
org.apache.pulsar.client.api.SubscriptionType.Shared;
- }
- }
- }
-
public enum Runtime {
JAVA,
PYTHON
@@ -98,7 +84,6 @@
private String logTopic;
private ProcessingGuarantees processingGuarantees;
private Map<String, Object> userConfig;
- private SubscriptionType subscriptionType;
private Runtime runtime;
private boolean autoAck;
@isPositiveNumber
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 64c23f537b..bc86c2d42a 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
@@ -450,13 +450,6 @@ private static void doCommonChecks(FunctionConfig
functionConfig) {
// Ensure that topics aren't being used as both input and output
verifyNoTopicClash(functionConfig.getInputs(),
functionConfig.getOutput());
- if (functionConfig.getSubscriptionType() != null
- && functionConfig.getSubscriptionType() !=
FunctionConfig.SubscriptionType.FAILOVER
- && functionConfig.getProcessingGuarantees() != null
- && functionConfig.getProcessingGuarantees() ==
FunctionConfig.ProcessingGuarantees.EFFECTIVELY_ONCE) {
- throw new IllegalArgumentException("Effectively-once
processing semantics can only be achieved using a Failover subscription type");
- }
-
WindowConfig windowConfig = functionConfig.getWindowConfig();
if (windowConfig != null) {
// set auto ack to false since windowing framework is
responsible
----------------------------------------------------------------
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