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)); + }); + } +}
