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 ce3d8d323685f..889d90e76e95d 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 @@ -297,10 +297,11 @@ abstract class FunctionDetailsCommand extends BaseCommand { @Option(names = "--retainOrdering", description = "Function consumes and processes messages in order", hidden = true) protected Boolean deprecatedRetainOrdering; - @Option(names = "--retain-ordering", description = "Function consumes and processes messages in order #Java") + @Option(names = "--retain-ordering", + description = "Function consumes and processes messages in order #Java, Python, Go") protected Boolean retainOrdering; @Option(names = "--retain-key-ordering", - description = "Function consumes and processes messages in key order #Java") + description = "Function consumes and processes messages in key order #Java, Python, Go") protected Boolean retainKeyOrdering; @Option(names = "--batch-builder", description = "BatcherBuilder provides two types of " + "batch construction methods, DEFAULT and KEY_BASED. The default value is: DEFAULT") diff --git a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfigUtils.java b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfigUtils.java index d88ac2820a966..81a3f50508cdd 100644 --- a/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfigUtils.java +++ b/pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfigUtils.java @@ -774,10 +774,6 @@ private static void doGolangChecks(FunctionConfig functionConfig) { if (functionConfig.getMaxMessageRetries() != null && functionConfig.getMaxMessageRetries() >= 0) { throw new IllegalArgumentException("Message retries not yet supported in Go function"); } - - if (functionConfig.getRetainKeyOrdering() != null && functionConfig.getRetainKeyOrdering()) { - throw new IllegalArgumentException("Retain Key Orderering not yet supported in Go function"); - } } private static void verifyNoTopicClash(Collection inputTopics, String outputTopic) diff --git a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/FunctionConfigUtilsTest.java b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/FunctionConfigUtilsTest.java index ba5e40429cf2d..dd0eb8d98b709 100644 --- a/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/FunctionConfigUtilsTest.java +++ b/pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/FunctionConfigUtilsTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.functions.utils; import static org.apache.pulsar.common.functions.FunctionConfig.ProcessingGuarantees.EFFECTIVELY_ONCE; +import static org.apache.pulsar.common.functions.FunctionConfig.Runtime.GO; import static org.apache.pulsar.common.functions.FunctionConfig.Runtime.PYTHON; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; @@ -29,6 +30,7 @@ import java.lang.reflect.Field; import java.util.Arrays; import java.util.Collection; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.concurrent.atomic.AtomicReference; @@ -763,4 +765,69 @@ public void testConvertProducerSpecToProducerConfigAndBackToProducerSpec() { producerSpec.getCryptoSpec().getProducerEncryptionKeyNameAt(i)); } } + + private static FunctionConfig minimalGoFunctionConfig() { + FunctionConfig functionConfig = new FunctionConfig(); + functionConfig.setTenant("test-tenant"); + functionConfig.setNamespace("test-namespace"); + functionConfig.setName("test-function"); + functionConfig.setInputs(Collections.singletonList("persistent://public/default/input")); + functionConfig.setRuntime(GO); + functionConfig.setGo("/path/to/function"); + return functionConfig; + } + + @Test + public void testGoFunctionAcceptsRetainKeyOrdering() { + FunctionConfig functionConfig = minimalGoFunctionConfig(); + functionConfig.setRetainKeyOrdering(true); + + FunctionConfigUtils.validateNonJavaFunction(functionConfig); + + // The KeyShared subscription the Go runtime selects has to survive conversion, otherwise the + // instance never sees it. + assertEquals(FunctionConfigUtils.convert(functionConfig).getSource().getSubscriptionType(), + SubscriptionType.KEY_SHARED); + } + + @Test + public void testGoFunctionAcceptsRetainOrdering() { + FunctionConfig functionConfig = minimalGoFunctionConfig(); + functionConfig.setRetainOrdering(true); + + FunctionConfigUtils.validateNonJavaFunction(functionConfig); + + assertEquals(FunctionConfigUtils.convert(functionConfig).getSource().getSubscriptionType(), + SubscriptionType.FAILOVER); + } + + @Test(expectedExceptions = IllegalArgumentException.class, + expectedExceptionsMessageRegExp = "Only one of retain ordering or retain key ordering can be set") + public void testGoFunctionRejectsBothOrderingModes() { + FunctionConfig functionConfig = minimalGoFunctionConfig(); + functionConfig.setRetainOrdering(true); + functionConfig.setRetainKeyOrdering(true); + + FunctionConfigUtils.validateNonJavaFunction(functionConfig); + } + + @Test(expectedExceptions = IllegalArgumentException.class, + expectedExceptionsMessageRegExp = + "When effectively once processing guarantee is specified, retain Key ordering cannot be set") + public void testGoFunctionRejectsRetainKeyOrderingWithEffectivelyOnce() { + FunctionConfig functionConfig = minimalGoFunctionConfig(); + functionConfig.setRetainKeyOrdering(true); + functionConfig.setProcessingGuarantees(EFFECTIVELY_ONCE); + + FunctionConfigUtils.validateNonJavaFunction(functionConfig); + } + + @Test(expectedExceptions = IllegalArgumentException.class, + expectedExceptionsMessageRegExp = "Message retries not yet supported in Go function") + public void testGoFunctionStillRejectsMessageRetries() { + FunctionConfig functionConfig = minimalGoFunctionConfig(); + functionConfig.setMaxMessageRetries(3); + + FunctionConfigUtils.validateNonJavaFunction(functionConfig); + } }