Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> inputTopics, String outputTopic)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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);
}
}
Loading