From 083e9561669a44dc2f687e789f6eba7a7a23e94a Mon Sep 17 00:00:00 2001 From: Nikolay Degterinsky <43110995+evillique@users.noreply.github.com> Date: Thu, 17 Sep 2026 15:06:53 +0000 Subject: [PATCH] Merge pull request #113919 from kalavt/fix-kafka-consumers-with-assignment-metric Fix `KafkaConsumersWithAssignment` drifting negative Signed-off-by: Ilya Golshtein --- src/Storages/Kafka/KafkaConsumer.cpp | 9 ++-- .../test_storage_kafka/test_batch_fast.py | 47 +++++++++++++++++++ 2 files changed, 51 insertions(+), 5 deletions(-) diff --git a/src/Storages/Kafka/KafkaConsumer.cpp b/src/Storages/Kafka/KafkaConsumer.cpp index 7b7dcc2b68ef..89ae414e405f 100644 --- a/src/Storages/Kafka/KafkaConsumer.cpp +++ b/src/Storages/Kafka/KafkaConsumer.cpp @@ -107,11 +107,6 @@ void KafkaConsumer::createConsumer(cppkafka::Configuration consumer_config) // with topics/partitions we were working with before rebalance LOG_TRACE(log, "Rebalance initiated. Revoking partitions: {}", topic_partitions); - if (!topic_partitions.empty()) - { - CurrentMetrics::sub(CurrentMetrics::KafkaConsumersWithAssignment, 1); - } - // we can not flush data to target from that point (it is pulled, not pushed) // so the best we can now it to // 1) repeat last commit in sync mode (async could be still in queue, we need to be sure is is properly committed before rebalance) @@ -359,6 +354,10 @@ void KafkaConsumer::cleanUnprocessed() offsets_stored = 0; } +/// Sole owner of the decrements for `KafkaAssignedPartitions` and `KafkaConsumersWithAssignment`. +/// Every path that drops an assignment reaches it - the revocation callback, `moveConsumer` and +/// `KafkaConsumer`'s destructor - so a caller that also decrements on its own makes the gauge drift +/// by one per rebalance until it passes zero. void KafkaConsumer::cleanAssignment() { if (assignment.has_value()) diff --git a/tests/integration/test_storage_kafka/test_batch_fast.py b/tests/integration/test_storage_kafka/test_batch_fast.py index 8ecd5a7d49d3..1540123bb365 100644 --- a/tests/integration/test_storage_kafka/test_batch_fast.py +++ b/tests/integration/test_storage_kafka/test_batch_fast.py @@ -3635,6 +3635,53 @@ def test_kafka_assigned_partitions(kafka_cluster): ) +def test_kafka_consumers_with_assignment_after_rebalance(kafka_cluster): + suffix = k.random_string(6) + topic_name = f"consumers_with_assignment_{suffix}" + k.kafka_create_topic(k.get_admin_client(kafka_cluster), topic_name, num_partitions=2) + + metric_query = ( + "SELECT value FROM system.metrics WHERE metric = 'KafkaConsumersWithAssignment'" + ) + before = int(instance.query(metric_query)) + + def create(table): + instance.query( + f""" + CREATE TABLE test.{table} (key UInt64, value UInt64) + ENGINE = Kafka + SETTINGS kafka_broker_list = 'kafka1:19092', + kafka_topic_list = '{topic_name}', + kafka_group_name = '{topic_name}', + kafka_format = 'JSONEachRow'; + CREATE MATERIALIZED VIEW test.{table}_mv ENGINE = Memory AS SELECT * FROM test.{table}; + """ + ) + + # The first consumer takes both partitions. + create(f"kafka_a_{suffix}") + assert_eq_with_retry(instance, metric_query, str(before + 1)) + + # The second member joining the group revokes and reassigns the live assignment. + create(f"kafka_b_{suffix}") + assert_eq_with_retry( + instance, + f""" + SELECT num_rebalance_assignments + FROM system.kafka_consumers + WHERE database = 'test' AND table = 'kafka_b_{suffix}' + """, + "1", + ) + + for table in (f"kafka_a_{suffix}", f"kafka_b_{suffix}"): + instance.query(f"DROP TABLE test.{table}_mv SYNC") + instance.query(f"DROP TABLE test.{table} SYNC") + + # Without the fix the revocation is counted twice, so the gauge ends one below where it started. + assert_eq_with_retry(instance, metric_query, str(before)) + + if __name__ == "__main__": cluster.start() input("Cluster created, press any key to destroy...")