diff --git a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/InstanceConfig.java b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/InstanceConfig.java index c70c103402893..ddf437c192463 100644 --- a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/InstanceConfig.java +++ b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/InstanceConfig.java @@ -43,6 +43,7 @@ public class InstanceConfig { // Whether the pulsar admin client exposed to function context, default is disabled. @Getter private boolean exposePulsarAdminClientEnabled = false; + private int metricsPort; /** * Get the string representation of {@link #getInstanceId()}. @@ -56,4 +57,8 @@ public String getInstanceName() { public FunctionDetails getFunctionDetails() { return functionDetails; } + + public boolean hasValidMetricsPort() { + return metricsPort > 0 && metricsPort < 65536; + } } diff --git a/pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java b/pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java index 5c45ad9918161..61be4fa9591f3 100644 --- a/pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java +++ b/pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java @@ -417,6 +417,7 @@ private void startProcessMode(org.apache.pulsar.functions.proto.Function.Functio instanceConfig.setInstanceId(i + instanceIdOffset); instanceConfig.setMaxBufferedTuples(1024); instanceConfig.setPort(FunctionCommon.findAvailablePort()); + instanceConfig.setMetricsPort(FunctionCommon.findAvailablePort()); instanceConfig.setClusterName("local"); if (functionConfig != null) { instanceConfig.setMaxPendingAsyncRequests(functionConfig.getMaxPendingAsyncRequests()); @@ -499,6 +500,7 @@ private void startThreadedMode(org.apache.pulsar.functions.proto.Function.Functi instanceConfig.setInstanceId(i + instanceIdOffset); instanceConfig.setMaxBufferedTuples(1024); instanceConfig.setPort(FunctionCommon.findAvailablePort()); + instanceConfig.setMetricsPort(FunctionCommon.findAvailablePort()); instanceConfig.setClusterName("local"); if (functionConfig != null) { instanceConfig.setMaxPendingAsyncRequests(functionConfig.getMaxPendingAsyncRequests()); diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceStarter.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceStarter.java index 77441980c60d6..206711076ad15 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceStarter.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/JavaInstanceStarter.java @@ -173,6 +173,7 @@ public void start(String[] args, ClassLoader functionInstanceClassLoader, ClassL Function.FunctionDetails functionDetails = functionDetailsBuilder.build(); instanceConfig.setFunctionDetails(functionDetails); instanceConfig.setPort(port); + instanceConfig.setMetricsPort(metrics_port); Map secretsProviderConfigMap = null; if (!StringUtils.isEmpty(secretsProviderConfig)) { diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/RuntimeUtils.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/RuntimeUtils.java index 0c7db6ef83a6e..4204d3a7877b7 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/RuntimeUtils.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/RuntimeUtils.java @@ -70,7 +70,6 @@ public static List composeCmd(InstanceConfig instanceConfig, Boolean installUserCodeDependencies, String pythonDependencyRepository, String pythonExtraDependencyRepository, - int metricsPort, String narExtractionDirectory, String functionInstanceClassPath, String pulsarWebServiceUrl) throws Exception { @@ -82,7 +81,7 @@ public static List composeCmd(InstanceConfig instanceConfig, authConfig, shardId, grpcPort, expectedHealthCheckInterval, logConfigFile, secretsProviderClassName, secretsProviderConfig, installUserCodeDependencies, pythonDependencyRepository, - pythonExtraDependencyRepository, metricsPort, narExtractionDirectory, + pythonExtraDependencyRepository, narExtractionDirectory, functionInstanceClassPath, false, pulsarWebServiceUrl)); return cmd; } @@ -120,8 +119,7 @@ public static List getArgsBeforeCmd(InstanceConfig instanceConfig, Strin public static List getGoInstanceCmd(InstanceConfig instanceConfig, String originalCodeFileName, String pulsarServiceUrl, - boolean k8sRuntime, - int metricsPort) throws IOException { + boolean k8sRuntime) throws IOException { final List args = new LinkedList<>(); GoInstanceConfig goInstanceConfig = new GoInstanceConfig(); @@ -221,8 +219,8 @@ public static List getGoInstanceCmd(InstanceConfig instanceConfig, goInstanceConfig.setMaxMessageRetries(instanceConfig.getFunctionDetails().getRetryDetails().getMaxMessageRetries()); } - if (metricsPort > 0 && metricsPort < 65536) { - goInstanceConfig.setMetricsPort(metricsPort); + if (instanceConfig.hasValidMetricsPort()) { + goInstanceConfig.setMetricsPort(instanceConfig.getMetricsPort()); } goInstanceConfig.setKillAfterIdleMs(0); @@ -260,7 +258,6 @@ public static List getCmd(InstanceConfig instanceConfig, Boolean installUserCodeDependencies, String pythonDependencyRepository, String pythonExtraDependencyRepository, - int metricsPort, String narExtractionDirectory, String functionInstanceClassPath, boolean k8sRuntime, @@ -268,7 +265,7 @@ public static List getCmd(InstanceConfig instanceConfig, final List args = new LinkedList<>(); if (instanceConfig.getFunctionDetails().getRuntime() == Function.FunctionDetails.Runtime.GO) { - return getGoInstanceCmd(instanceConfig, originalCodeFileName, pulsarServiceUrl, k8sRuntime, metricsPort); + return getGoInstanceCmd(instanceConfig, originalCodeFileName, pulsarServiceUrl, k8sRuntime); } if (instanceConfig.getFunctionDetails().getRuntime() == Function.FunctionDetails.Runtime.JAVA) { @@ -398,7 +395,7 @@ && isNotBlank(authConfig.getClientAuthenticationParameters())) { args.add(String.valueOf(grpcPort)); args.add("--metrics_port"); - args.add(String.valueOf(metricsPort)); + args.add(String.valueOf(instanceConfig.getMetricsPort())); // state storage configs if (null != stateStorageServiceUrl) { diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntime.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntime.java index 4ddeb0c81de2f..693d787c8fedc 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntime.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntime.java @@ -189,7 +189,6 @@ public class KubernetesRuntime implements Runtime { Optional functionAuthDataCacheProvider, boolean authenticationEnabled, Integer grpcPort, - Integer metricsPort, String narExtractionDirectory, Optional manifestCustomizer, String functinoInstanceClassPath, @@ -238,7 +237,7 @@ public class KubernetesRuntime implements Runtime { this.functionAuthDataCacheProvider = functionAuthDataCacheProvider; this.grpcPort = grpcPort; - this.metricsPort = metricsPort; + this.metricsPort = instanceConfig.hasValidMetricsPort() ? instanceConfig.getMetricsPort() : null; this.narExtractionDirectory = narExtractionDirectory; this.processArgs = new LinkedList<>(); @@ -275,7 +274,6 @@ public class KubernetesRuntime implements Runtime { installUserCodeDependencies, pythonDependencyRepository, pythonExtraDependencyRepository, - metricsPort, narExtractionDirectory, functinoInstanceClassPath, true, diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactory.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactory.java index 0155e2a2f35a3..9155cf3dcc14d 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactory.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeFactory.java @@ -317,7 +317,6 @@ public KubernetesRuntime createContainer(InstanceConfig instanceConfig, String c authProvider, authenticationEnabled, grpcPort, - metricsPort, narExtractionDirectory, manifestCustomizer, functionInstanceClassPath, diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/process/ProcessRuntime.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/process/ProcessRuntime.java index b6fcb0057e81d..f271b3326bfef 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/process/ProcessRuntime.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/process/ProcessRuntime.java @@ -91,7 +91,7 @@ class ProcessRuntime implements Runtime { String pulsarWebServiceUrl) throws Exception { this.instanceConfig = instanceConfig; this.instancePort = instanceConfig.getPort(); - this.metricsPort = FunctionCommon.findAvailablePort(); + this.metricsPort = instanceConfig.getMetricsPort(); this.expectedHealthCheckInterval = expectedHealthCheckInterval; this.secretsProviderConfigurator = secretsProviderConfigurator; this.funcLogDir = RuntimeUtils.genFunctionLogFolder(logDirectory, instanceConfig); @@ -134,7 +134,6 @@ class ProcessRuntime implements Runtime { false, null, null, - this.metricsPort, narExtractionDirectory, null, pulsarWebServiceUrl); diff --git a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/RuntimeUtilsTest.java b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/RuntimeUtilsTest.java index 72dd61b911f19..f8bbbc4a883af 100644 --- a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/RuntimeUtilsTest.java +++ b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/RuntimeUtilsTest.java @@ -65,6 +65,7 @@ public void getGoInstanceCmd(boolean k8sRuntime) throws IOException { instanceConfig.setFunctionVersion("1.0.0"); instanceConfig.setMaxBufferedTuples(5); instanceConfig.setPort(1337); + instanceConfig.setMetricsPort(60000); JSONObject userConfig = new JSONObject(); @@ -108,7 +109,7 @@ public void getGoInstanceCmd(boolean k8sRuntime) throws IOException { instanceConfig.setFunctionDetails(functionDetails); - List commands = RuntimeUtils.getGoInstanceCmd(instanceConfig, "config", "pulsar://localhost:6650", k8sRuntime, 60000); + List commands = RuntimeUtils.getGoInstanceCmd(instanceConfig, "config", "pulsar://localhost:6650", k8sRuntime); if (k8sRuntime) { goInstanceConfig = new ObjectMapper().readValue(commands.get(2).replaceAll("^\'|\'$", ""), HashMap.class); } else { diff --git a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java index cd621d21c84bf..6eadad947d28a 100644 --- a/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java +++ b/pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/kubernetes/KubernetesRuntimeTest.java @@ -876,6 +876,7 @@ InstanceConfig createGolangInstanceConfig() { config.setInstanceId(0); config.setMaxBufferedTuples(1024); config.setClusterName("standalone"); + config.setMetricsPort(4331); return config; } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionActioner.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionActioner.java index 73f8f9540350e..6cf90606d4848 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionActioner.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionActioner.java @@ -184,6 +184,7 @@ InstanceConfig createInstanceConfig(FunctionDetails functionDetails, Function.Fu instanceConfig.setInstanceId(instanceId); instanceConfig.setMaxBufferedTuples(1024); instanceConfig.setPort(FunctionCommon.findAvailablePort()); + instanceConfig.setMetricsPort(FunctionCommon.findAvailablePort()); instanceConfig.setClusterName(clusterName); instanceConfig.setFunctionAuthenticationSpec(functionAuthSpec); instanceConfig.setMaxPendingAsyncRequests(workerConfig.getMaxPendingAsyncRequests());