Skip to content

Commit ecf6dae

Browse files
[fix][fn] Allow retainKeyOrdering on Go functions (#26421)
1 parent 4040eec commit ecf6dae

3 files changed

Lines changed: 70 additions & 6 deletions

File tree

pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -297,10 +297,11 @@ abstract class FunctionDetailsCommand extends BaseCommand {
297297
@Option(names = "--retainOrdering",
298298
description = "Function consumes and processes messages in order", hidden = true)
299299
protected Boolean deprecatedRetainOrdering;
300-
@Option(names = "--retain-ordering", description = "Function consumes and processes messages in order #Java")
300+
@Option(names = "--retain-ordering",
301+
description = "Function consumes and processes messages in order #Java, Python, Go")
301302
protected Boolean retainOrdering;
302303
@Option(names = "--retain-key-ordering",
303-
description = "Function consumes and processes messages in key order #Java")
304+
description = "Function consumes and processes messages in key order #Java, Python, Go")
304305
protected Boolean retainKeyOrdering;
305306
@Option(names = "--batch-builder", description = "BatcherBuilder provides two types of "
306307
+ "batch construction methods, DEFAULT and KEY_BASED. The default value is: DEFAULT")

pulsar-functions/utils/src/main/java/org/apache/pulsar/functions/utils/FunctionConfigUtils.java

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -774,10 +774,6 @@ private static void doGolangChecks(FunctionConfig functionConfig) {
774774
if (functionConfig.getMaxMessageRetries() != null && functionConfig.getMaxMessageRetries() >= 0) {
775775
throw new IllegalArgumentException("Message retries not yet supported in Go function");
776776
}
777-
778-
if (functionConfig.getRetainKeyOrdering() != null && functionConfig.getRetainKeyOrdering()) {
779-
throw new IllegalArgumentException("Retain Key Orderering not yet supported in Go function");
780-
}
781777
}
782778

783779
private static void verifyNoTopicClash(Collection<String> inputTopics, String outputTopic)

pulsar-functions/utils/src/test/java/org/apache/pulsar/functions/utils/FunctionConfigUtilsTest.java

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
package org.apache.pulsar.functions.utils;
2020

2121
import static org.apache.pulsar.common.functions.FunctionConfig.ProcessingGuarantees.EFFECTIVELY_ONCE;
22+
import static org.apache.pulsar.common.functions.FunctionConfig.Runtime.GO;
2223
import static org.apache.pulsar.common.functions.FunctionConfig.Runtime.PYTHON;
2324
import static org.testng.Assert.assertEquals;
2425
import static org.testng.Assert.assertFalse;
@@ -29,6 +30,7 @@
2930
import java.lang.reflect.Field;
3031
import java.util.Arrays;
3132
import java.util.Collection;
33+
import java.util.Collections;
3234
import java.util.HashMap;
3335
import java.util.Map;
3436
import java.util.concurrent.atomic.AtomicReference;
@@ -763,4 +765,69 @@ public void testConvertProducerSpecToProducerConfigAndBackToProducerSpec() {
763765
producerSpec.getCryptoSpec().getProducerEncryptionKeyNameAt(i));
764766
}
765767
}
768+
769+
private static FunctionConfig minimalGoFunctionConfig() {
770+
FunctionConfig functionConfig = new FunctionConfig();
771+
functionConfig.setTenant("test-tenant");
772+
functionConfig.setNamespace("test-namespace");
773+
functionConfig.setName("test-function");
774+
functionConfig.setInputs(Collections.singletonList("persistent://public/default/input"));
775+
functionConfig.setRuntime(GO);
776+
functionConfig.setGo("/path/to/function");
777+
return functionConfig;
778+
}
779+
780+
@Test
781+
public void testGoFunctionAcceptsRetainKeyOrdering() {
782+
FunctionConfig functionConfig = minimalGoFunctionConfig();
783+
functionConfig.setRetainKeyOrdering(true);
784+
785+
FunctionConfigUtils.validateNonJavaFunction(functionConfig);
786+
787+
// The KeyShared subscription the Go runtime selects has to survive conversion, otherwise the
788+
// instance never sees it.
789+
assertEquals(FunctionConfigUtils.convert(functionConfig).getSource().getSubscriptionType(),
790+
SubscriptionType.KEY_SHARED);
791+
}
792+
793+
@Test
794+
public void testGoFunctionAcceptsRetainOrdering() {
795+
FunctionConfig functionConfig = minimalGoFunctionConfig();
796+
functionConfig.setRetainOrdering(true);
797+
798+
FunctionConfigUtils.validateNonJavaFunction(functionConfig);
799+
800+
assertEquals(FunctionConfigUtils.convert(functionConfig).getSource().getSubscriptionType(),
801+
SubscriptionType.FAILOVER);
802+
}
803+
804+
@Test(expectedExceptions = IllegalArgumentException.class,
805+
expectedExceptionsMessageRegExp = "Only one of retain ordering or retain key ordering can be set")
806+
public void testGoFunctionRejectsBothOrderingModes() {
807+
FunctionConfig functionConfig = minimalGoFunctionConfig();
808+
functionConfig.setRetainOrdering(true);
809+
functionConfig.setRetainKeyOrdering(true);
810+
811+
FunctionConfigUtils.validateNonJavaFunction(functionConfig);
812+
}
813+
814+
@Test(expectedExceptions = IllegalArgumentException.class,
815+
expectedExceptionsMessageRegExp =
816+
"When effectively once processing guarantee is specified, retain Key ordering cannot be set")
817+
public void testGoFunctionRejectsRetainKeyOrderingWithEffectivelyOnce() {
818+
FunctionConfig functionConfig = minimalGoFunctionConfig();
819+
functionConfig.setRetainKeyOrdering(true);
820+
functionConfig.setProcessingGuarantees(EFFECTIVELY_ONCE);
821+
822+
FunctionConfigUtils.validateNonJavaFunction(functionConfig);
823+
}
824+
825+
@Test(expectedExceptions = IllegalArgumentException.class,
826+
expectedExceptionsMessageRegExp = "Message retries not yet supported in Go function")
827+
public void testGoFunctionStillRejectsMessageRetries() {
828+
FunctionConfig functionConfig = minimalGoFunctionConfig();
829+
functionConfig.setMaxMessageRetries(3);
830+
831+
FunctionConfigUtils.validateNonJavaFunction(functionConfig);
832+
}
766833
}

0 commit comments

Comments
 (0)