This is an automated email from the ASF dual-hosted git repository.

gyfora pushed a commit to branch release-1.16
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git

commit aea8ceb0ee1d95e0c2c8793100c295b4a665bd4d
Author: Purushottam Sinha <[email protected]>
AuthorDate: Thu Aug 20 20:37:03 2026 +0530

    [FLINK-40443] Harden autoscaler metrics name parsing against complex regex 
(#1187)
---
 .../flink/autoscaler/ScalingMetricCollector.java   |  34 +--
 .../utils/PartitionMetricNameParser.java           | 105 ++++++++
 .../utils/PartitionMetricNameParserTest.java       | 278 +++++++++++++++++++++
 3 files changed, 392 insertions(+), 25 deletions(-)

diff --git 
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java
 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java
index 9cdf4ace..578604d2 100644
--- 
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java
+++ 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java
@@ -31,6 +31,7 @@ import org.apache.flink.autoscaler.metrics.ScalingMetrics;
 import org.apache.flink.autoscaler.state.AutoScalerStateStore;
 import org.apache.flink.autoscaler.topology.IOMetrics;
 import org.apache.flink.autoscaler.topology.JobTopology;
+import org.apache.flink.autoscaler.utils.PartitionMetricNameParser;
 import org.apache.flink.client.program.rest.RestClusterClient;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.runtime.execution.ExecutionState;
@@ -65,8 +66,6 @@ import java.util.SortedMap;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.TimeUnit;
 import java.util.function.Supplier;
-import java.util.regex.Matcher;
-import java.util.regex.Pattern;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
@@ -306,35 +305,20 @@ public abstract class ScalingMetricCollector<KEY, Context 
extends JobAutoScalerC
     private void updateKafkaPulsarSourceNumPartitions(
             Context ctx, JobID jobId, JobTopology topology) throws Exception {
         try (var restClient = ctx.getRestClusterClient()) {
-            Pattern partitionRegex =
-                    Pattern.compile(
-                            
"^.*?(\\.kafkaCluster\\.(?<kafkaCluster>.+))?\\.KafkaSourceReader\\.topic\\.(?<kafkaTopic>.+)\\.partition\\.(?<kafkaId>\\d+)\\.currentOffset$"
-                                    + 
"|^.*\\.PulsarConsumer\\.(?<pulsarTopic>.+)-partition-(?<pulsarId>\\d+)\\..*\\.numMsgsReceived$");
             for (var vertexInfo : topology.getVertexInfos().values()) {
-                if (vertexInfo.getInputs().isEmpty()) {
+                if (topology.isSource(vertexInfo.getId())) {
                     var sourceVertex = vertexInfo.getId();
                     var numPartitions =
                             queryAggregatedMetricNames(restClient, jobId, 
sourceVertex).stream()
                                     .map(
                                             v -> {
-                                                Matcher matcher = 
partitionRegex.matcher(v);
-                                                if (matcher.matches()) {
-                                                    String kafkaTopic = 
matcher.group("kafkaTopic");
-                                                    String kafkaCluster =
-                                                            
matcher.group("kafkaCluster");
-                                                    String kafkaId = 
matcher.group("kafkaId");
-                                                    String pulsarTopic =
-                                                            
matcher.group("pulsarTopic");
-                                                    String pulsarId = 
matcher.group("pulsarId");
-                                                    return kafkaTopic != null
-                                                            ? kafkaCluster
-                                                                    + "-"
-                                                                    + 
kafkaTopic
-                                                                    + "-"
-                                                                    + kafkaId
-                                                            : pulsarTopic + 
"-" + pulsarId;
-                                                }
-                                                return null;
+                                                String key =
+                                                        
PartitionMetricNameParser
+                                                                
.parseKafkaPartitionKey(v);
+                                                return key != null
+                                                        ? key
+                                                        : 
PartitionMetricNameParser
+                                                                
.parsePulsarPartitionKey(v);
                                             })
                                     .filter(Objects::nonNull)
                                     .collect(Collectors.toSet())
diff --git 
a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParser.java
 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParser.java
new file mode 100644
index 00000000..cccd4de3
--- /dev/null
+++ 
b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParser.java
@@ -0,0 +1,105 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.autoscaler.utils;
+
+/**
+ * Parses Kafka/Pulsar source-partition metric names to extract a stable 
per-partition key.
+ *
+ * <p>Uses plain string splitting/indexing rather than a regex, since metric 
names are fully
+ * job-controlled and this must stay linear-time regardless of input shape.
+ *
+ * <p>Parsing splits on the '.' delimiter and treats each name component as 
dot-free: Flink's metric
+ * query service replaces '.' within a component with '_' before joining, so a 
'.' in the
+ * identifiers returned by the REST metrics endpoint is always a real 
delimiter. Only correctness on
+ * legitimate input relies on this assumption; safety does not, since an 
unexpected shape merely
+ * yields null.
+ */
+public final class PartitionMetricNameParser {
+
+    private static final String KAFKA_CLUSTER = "kafkaCluster";
+    private static final String KAFKA_SOURCE_READER = "KafkaSourceReader";
+    private static final String TOPIC = "topic";
+    private static final String PARTITION = "partition";
+    private static final String CURRENT_OFFSET = "currentOffset";
+    private static final String PULSAR_CONSUMER = "PulsarConsumer";
+    private static final String NUM_MSGS_RECEIVED = "numMsgsReceived";
+    private static final String PULSAR_PARTITION_INFIX = "-partition-";
+
+    private PartitionMetricNameParser() {}
+
+    /**
+     * Returns {@code "<cluster>-<topic>-<id>"}, or {@code null} if this is 
not a Kafka
+     * currentOffset partition metric. The cluster segment is the literal 
"null" when absent; only
+     * distinct keys matter to callers.
+     */
+    public static String parseKafkaPartitionKey(String metricName) {
+        String[] parts = metricName.split("\\.", -1);
+        int n = parts.length;
+        // tail: KafkaSourceReader.topic.<topic>.partition.<id>.currentOffset
+        if (n < 6
+                || !parts[n - 6].equals(KAFKA_SOURCE_READER)
+                || !parts[n - 5].equals(TOPIC)
+                || !parts[n - 3].equals(PARTITION)
+                || !isAllDigits(parts[n - 2])
+                || !parts[n - 1].equals(CURRENT_OFFSET)) {
+            return null;
+        }
+        // optional ".kafkaCluster.<cluster>." precedes the tail; null when 
absent
+        String cluster = (n >= 8 && parts[n - 8].equals(KAFKA_CLUSTER)) ? 
parts[n - 7] : null;
+        return cluster + "-" + parts[n - 4] + "-" + parts[n - 2];
+    }
+
+    /**
+     * Returns {@code "<topic>-<id>"}, or {@code null} if this is not a Pulsar 
numMsgsReceived
+     * partition metric.
+     */
+    public static String parsePulsarPartitionKey(String metricName) {
+        String[] parts = metricName.split("\\.", -1);
+        int n = parts.length;
+        // tail: 
PulsarConsumer.<topic>-partition-<id>.<consumerHash>.numMsgsReceived
+        if (n < 3 || !parts[n - 1].equals(NUM_MSGS_RECEIVED)) {
+            return null;
+        }
+        for (int i = 0; i + 1 < n - 1; i++) {
+            if (!parts[i].equals(PULSAR_CONSUMER)) {
+                continue;
+            }
+            String segment = parts[i + 1]; // "<topic>-partition-<id>"
+            int idx = segment.lastIndexOf(PULSAR_PARTITION_INFIX);
+            if (idx <= 0) {
+                return null;
+            }
+            String id = segment.substring(idx + 
PULSAR_PARTITION_INFIX.length());
+            return isAllDigits(id) ? segment.substring(0, idx) + "-" + id : 
null;
+        }
+        return null;
+    }
+
+    private static boolean isAllDigits(String s) {
+        if (s.isEmpty()) {
+            return false; // matches \d+, which requires at least one digit
+        }
+        for (int i = 0; i < s.length(); i++) {
+            char c = s.charAt(i);
+            if (c < '0' || c > '9') { // ASCII only, matching regex \d
+                return false;
+            }
+        }
+        return true;
+    }
+}
diff --git 
a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParserTest.java
 
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParserTest.java
new file mode 100644
index 00000000..685e8ddc
--- /dev/null
+++ 
b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/utils/PartitionMetricNameParserTest.java
@@ -0,0 +1,278 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.autoscaler.utils;
+
+import org.junit.jupiter.api.Test;
+
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Random;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+
+/** Test for {@link PartitionMetricNameParser}. */
+public class PartitionMetricNameParserTest {
+
+    /**
+     * The original regex-based implementation that {@link 
PartitionMetricNameParser} replaced, kept
+     * here so we can assert the new parser produces exactly the same result 
on realistic inputs
+     * (see {@link #testEquivalentToLegacyRegex()}).
+     */
+    private static final Pattern LEGACY_REGEX =
+            Pattern.compile(
+                    
"^.*?(\\.kafkaCluster\\.(?<kafkaCluster>.+))?\\.KafkaSourceReader\\.topic\\.(?<kafkaTopic>.+)\\.partition\\.(?<kafkaId>\\d+)\\.currentOffset$"
+                            + 
"|^.*\\.PulsarConsumer\\.(?<pulsarTopic>.+)-partition-(?<pulsarId>\\d+)\\..*\\.numMsgsReceived$");
+
+    private static String legacyPartitionKey(String metricName) {
+        Matcher matcher = LEGACY_REGEX.matcher(metricName);
+        if (matcher.matches()) {
+            String kafkaTopic = matcher.group("kafkaTopic");
+            String kafkaCluster = matcher.group("kafkaCluster");
+            String kafkaId = matcher.group("kafkaId");
+            String pulsarTopic = matcher.group("pulsarTopic");
+            String pulsarId = matcher.group("pulsarId");
+            return kafkaTopic != null
+                    ? kafkaCluster + "-" + kafkaTopic + "-" + kafkaId
+                    : pulsarTopic + "-" + pulsarId;
+        }
+        return null;
+    }
+
+    /** Mirrors how {@code ScalingMetricCollector} combines the two parser 
methods. */
+    private static String newPartitionKey(String metricName) {
+        String kafka = 
PartitionMetricNameParser.parseKafkaPartitionKey(metricName);
+        return kafka != null
+                ? kafka
+                : 
PartitionMetricNameParser.parsePulsarPartitionKey(metricName);
+    }
+
+    @Test
+    public void testKafkaWithoutCluster() {
+        assertEquals(
+                "null-testTopic-0",
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.testTopic.partition.0.currentOffset"));
+        assertEquals(
+                "null-anotherTopic-0",
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.anotherTopic.partition.0.currentOffset"));
+    }
+
+    @Test
+    public void testKafkaWithCluster() {
+        assertEquals(
+                "my-cluster-1-testTopic-0",
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"1.Source__Kafka_Source_(testTopic).kafkaCluster.my-cluster-1.KafkaSourceReader.topic.testTopic.partition.0.currentOffset"));
+        assertEquals(
+                "my-cluster-2-testTopic-3",
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"1.Source__Kafka_Source_(testTopic).kafkaCluster.my-cluster-2.KafkaSourceReader.topic.testTopic.partition.3.currentOffset"));
+    }
+
+    @Test
+    public void testKafkaNonMatchingMetric() {
+        assertNull(
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.testTopic.partition.0.anotherMetric"));
+    }
+
+    @Test
+    public void testKafkaMalformedShapes() {
+        // partition id is not numeric
+        assertNull(
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"x.KafkaSourceReader.topic.testTopic.partition.abc.currentOffset"));
+        // non-ASCII digits do not count as \d
+        assertNull(
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"x.KafkaSourceReader.topic.testTopic.partition.٥.currentOffset"));
+        // empty partition id segment
+        assertNull(
+                PartitionMetricNameParser.parseKafkaPartitionKey(
+                        
"x.KafkaSourceReader.topic.testTopic.partition..currentOffset"));
+        // too short to contain the required shape
+        
assertNull(PartitionMetricNameParser.parseKafkaPartitionKey("topic.testTopic.partition"));
+        assertNull(PartitionMetricNameParser.parseKafkaPartitionKey(""));
+    }
+
+    @Test
+    public void testPulsar() {
+        assertEquals(
+                "persistent_//public/default/testTopic-1",
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/testTopic-partition-1.d842f.numMsgsReceived"));
+        // Same topic/partition, different (irrelevant) consumer-hash segment 
-> same key.
+        assertEquals(
+                "persistent_//public/default/testTopic-1",
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/testTopic-partition-1.660d2.numMsgsReceived"));
+        assertEquals(
+                "persistent_//public/default/otherTopic-2",
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/otherTopic-partition-2.m953d.numMsgsReceived"));
+    }
+
+    @Test
+    public void testPulsarNonMatchingMetric() {
+        assertNull(
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/testTopic-partition-1.d842f.someOtherMetric"));
+    }
+
+    @Test
+    public void testPulsarMalformedShapes() {
+        // partition id is not numeric
+        assertNull(
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.PulsarConsumer.testTopic-partition-abc.numMsgsReceived"));
+        // no "-partition-" infix at all
+        assertNull(
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        "0.PulsarConsumer.testTopic.numMsgsReceived"));
+        // too short to contain the required shape
+        
assertNull(PartitionMetricNameParser.parsePulsarPartitionKey("PulsarConsumer"));
+        assertNull(PartitionMetricNameParser.parsePulsarPartitionKey(""));
+    }
+
+    @Test
+    public void testPulsarFirstMarkerWinsWithoutException() {
+        // only the first "PulsarConsumer" segment is considered
+        assertNull(
+                PartitionMetricNameParser.parsePulsarPartitionKey(
+                        
"0.PulsarConsumer.noPartitionHere.PulsarConsumer.testTopic-partition-1.numMsgsReceived"));
+    }
+
+    /**
+     * Asserts the new parser produces exactly the same key as the original 
regex on the captured
+     * real-world metric names plus a large set of generated realistic ones. 
Inputs are kept within
+     * the shape actually emitted by Flink's metric query service (each 
metric-name component is
+     * dot-free, since '.' is the delimiter), which is the space the operator 
ever sees in practice.
+     */
+    @Test
+    public void testEquivalentToLegacyRegex() {
+        List<String> inputs = new ArrayList<>();
+
+        // Captured real-world fixtures (mirroring 
MetricsCollectionAndEvaluationTest).
+        inputs.add(
+                
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.testTopic.partition.0.anotherMetric");
+        inputs.add(
+                
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.anotherTopic.partition.0.currentOffset");
+        inputs.add(
+                
"1.Source__Kafka_Source_(testTopic).KafkaSourceReader.topic.testTopic.partition.3.currentOffset");
+        inputs.add(
+                
"1.Source__Kafka_Source_(testTopic).kafkaCluster.my-cluster-2.KafkaSourceReader.topic.testTopic.partition.1.currentOffset");
+        inputs.add(
+                
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/testTopic-partition-1.d842f.numMsgsReceived");
+        inputs.add(
+                
"0.Source__pulsar_source[1].PulsarConsumer.persistent_//public/default/otherTopic-partition-2.m953d.numMsgsReceived");
+
+        // Generated realistic names (deterministic seed for reproducibility). 
Kafka and Pulsar use
+        // connector-specific source-operator prefixes and topic shapes, 
mirroring what each emits.
+        Random rnd = new Random(42);
+        String[] kafkaTopics = {"testTopic", "orders_eu_v1", "anotherTopic", 
"my-topic", "a"};
+        String[] pulsarTopics = {
+            "persistent_//public/default/testTopic",
+            "persistent_//public/default/orders_eu_v1",
+            "persistent_//tenant/ns/my-topic",
+            "non-persistent_//public/default/events"
+        };
+        String[] clusters = {"my-cluster-1", "cluster_2", "c"};
+        for (int i = 0; i < 100; i++) {
+            int id = rnd.nextInt(1000);
+            String hash = Integer.toHexString(rnd.nextInt(0xfffff));
+
+            String kafkaTopic = kafkaTopics[rnd.nextInt(kafkaTopics.length)];
+            String kafkaPrefix = rnd.nextInt(9) + ".Source__Kafka_Source_(" + 
kafkaTopic + ")";
+            // Kafka without cluster.
+            inputs.add(
+                    kafkaPrefix
+                            + ".KafkaSourceReader.topic."
+                            + kafkaTopic
+                            + ".partition."
+                            + id
+                            + ".currentOffset");
+            // Kafka with cluster.
+            inputs.add(
+                    kafkaPrefix
+                            + ".kafkaCluster."
+                            + clusters[rnd.nextInt(clusters.length)]
+                            + ".KafkaSourceReader.topic."
+                            + kafkaTopic
+                            + ".partition."
+                            + id
+                            + ".currentOffset");
+            // Kafka near-miss negative (wrong suffix).
+            inputs.add(
+                    kafkaPrefix
+                            + ".KafkaSourceReader.topic."
+                            + kafkaTopic
+                            + ".partition."
+                            + id
+                            + ".someOtherMetric");
+
+            String pulsarTopic = 
pulsarTopics[rnd.nextInt(pulsarTopics.length)];
+            String pulsarPrefix = rnd.nextInt(9) + ".Source__pulsar_source[" + 
rnd.nextInt(4) + "]";
+            // Pulsar.
+            inputs.add(
+                    pulsarPrefix
+                            + ".PulsarConsumer."
+                            + pulsarTopic
+                            + "-partition-"
+                            + id
+                            + "."
+                            + hash
+                            + ".numMsgsReceived");
+            // Pulsar near-miss negative (non-numeric id).
+            inputs.add(
+                    pulsarPrefix
+                            + ".PulsarConsumer."
+                            + pulsarTopic
+                            + "-partition-x."
+                            + hash
+                            + ".numMsgsReceived");
+        }
+
+        for (String v : inputs) {
+            assertEquals(legacyPartitionKey(v), newPartitionKey(v), "mismatch 
for: " + v);
+        }
+    }
+
+    @Test
+    public void testAdversarialPayloadDoesNotHang() {
+        StringBuilder evil = new StringBuilder("0.Source__pulsar_source[1]");
+        for (int i = 0; i < 2000; i++) {
+            evil.append(".PulsarConsumer.X-partition-1.a.numMsgsReceivedZ");
+        }
+        evil.append(".end");
+        String payload = evil.toString();
+
+        assertTimeoutPreemptively(
+                Duration.ofSeconds(5),
+                () -> {
+                    
assertNull(PartitionMetricNameParser.parseKafkaPartitionKey(payload));
+                    
assertNull(PartitionMetricNameParser.parsePulsarPartitionKey(payload));
+                });
+    }
+}

Reply via email to