Skip to content
Open
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
9 changes: 4 additions & 5 deletions src/Storages/Kafka/KafkaConsumer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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())
Expand Down
47 changes: 47 additions & 0 deletions tests/integration/test_storage_kafka/test_batch_fast.py
Original file line number Diff line number Diff line change
Expand Up @@ -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...")
Expand Down
Loading