Skip to content

perf(assignor): optimize cooperative-sticky assignment - #5551

Open
Kyle Phelps (kphelps) wants to merge 8 commits into
confluentinc:masterfrom
DataDog:kphelps/cooperative-sticky-optimizations
Open

perf(assignor): optimize cooperative-sticky assignment#5551
Kyle Phelps (kphelps) wants to merge 8 commits into
confluentinc:masterfrom
DataDog:kphelps/cooperative-sticky-optimizations

Conversation

@kphelps

@kphelps Kyle Phelps (kphelps) commented Jul 17, 2026

Copy link
Copy Markdown
Contributor

Problem

The cooperative-sticky assignor does substantially more work and allocation than necessary for large groups. I've seen assignments taking 10s of seconds, with the change that dropped to <1s.

Approach

This change keeps the cooperative-sticky algorithm and assignment behavior while reducing work in its hot paths:

  • detect fresh assignments from ownership state so cold runs take the initialization path
  • remove the redundant subscription comparison
  • size assignment, potential-partition, movement, and ownership maps from known counts
  • represent potential membership with eligible-topic references and counts rather than expanded consumer-to-partition entries
  • maintain changed-consumer ordering incrementally instead of repeatedly running qsort
  • skip rollback snapshots for cold assignments.

The changes are split into small commits so each optimization can be reviewed independently.

Outcome

Direct cooperative-sticky assignor benchmark: 1,000 members, 3,000 partitions, five cold runs and five warm runs, release builds with -O2 on an Apple M4 Max. This isolates assignor execution and excludes broker, network, and group-protocol overhead.

Revision Cold mean Warm mean Max RSS
master (4c3b017c) 12,239.5 ± 84.5 ms 12,372.6 ± 84.7 ms 619,675,648 B
This PR (d27694a7) 53.2 ± 2.9 ms 104.2 ± 1.7 ms 45,449,216 B
Improvement 229.9× faster 118.8× faster 92.7% lower
Benchmark code
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdint.h>

#include "rdkafka_int.h"
#include "rdkafka_assignor.h"
#include "rdkafka_metadata.h"

static int run_once(rd_kafka_t *rk,
                    const rd_kafka_assignor_t *rkas,
                    rd_kafka_metadata_t *metadata,
                    rd_kafka_group_member_t *members,
                    int member_cnt) {
        char errstr[256];
        for (int i = 0; i < member_cnt; i++)
                rd_list_clear(&members[i].rkgm_eligible);
        rd_kafka_resp_err_t err = rd_kafka_assignor_run(
            rk->rk_cgrp, rkas, metadata, members, member_cnt, errstr,
            sizeof(errstr));
        if (err) {
                fprintf(stderr, "assignor_run failed: %s (%s)\n", errstr,
                        rd_kafka_err2str(err));
                return 1;
        }
        return 0;
}

int main(int argc, char **argv) {
        const int member_cnt = argc > 1 ? atoi(argv[1]) : 1000;
        const int partition_cnt = argc > 2 ? atoi(argv[2]) : 3000;
        const int repetitions = argc > 3 ? atoi(argv[3]) : 5;
        rd_kafka_conf_t *conf = rd_kafka_conf_new();
        rd_kafka_t *rk;
        rd_kafka_assignor_t *rkas;
        rd_kafka_metadata_t *metadata;
        rd_kafka_group_member_t *members;
        char errstr[256];
        int i, r;

        if (rd_kafka_conf_set(conf, "group.id", "assignor-bench", errstr,
                              sizeof(errstr)) != RD_KAFKA_CONF_OK ||
            rd_kafka_conf_set(conf, "partition.assignment.strategy",
                              "cooperative-sticky", errstr,
                              sizeof(errstr)) != RD_KAFKA_CONF_OK) {
                fprintf(stderr, "conf failed: %s\n", errstr);
                return 1;
        }

        rk = rd_kafka_new(RD_KAFKA_CONSUMER, conf, errstr, sizeof(errstr));
        if (!rk) {
                fprintf(stderr, "client creation failed: %s\n", errstr);
                return 1;
        }

        rkas = rd_kafka_assignor_find(rk, "cooperative-sticky");
        if (!rkas) {
                fprintf(stderr, "cooperative-sticky assignor not found\n");
                rd_kafka_destroy(rk);
                return 1;
        }
        rd_kafka_cgrp_set_member_id(rk->rk_cgrp, "leader");

        metadata = rd_kafka_metadata_new_topic_mockv(1, "topic", partition_cnt);
        members = calloc((size_t)member_cnt, sizeof(*members));
        if (!metadata || !members) {
                fprintf(stderr, "allocation failed\n");
                return 1;
        }

        for (i = 0; i < member_cnt; i++) {
                char member_id[32];
                snprintf(member_id, sizeof(member_id), "member-%d", i);
                ut_init_member(&members[i], member_id, "topic", NULL);
        }

        /* Cold assignment: no member owns partitions yet. */
        for (r = 0; r < repetitions; r++) {
                uint64_t start = rd_clock();
                if (run_once(rk, rkas, metadata, members, member_cnt))
                        return 1;
                printf("cold[%d] %.3f ms\n", r,
                       (double)(rd_clock() - start) / 1000.0);
        }

        /* Warm assignment: report the generated assignment as owned. */
        for (i = 0; i < member_cnt; i++)
                ut_set_owned(&members[i]);
        for (r = 0; r < repetitions; r++) {
                uint64_t start = rd_clock();
                if (run_once(rk, rkas, metadata, members, member_cnt))
                        return 1;
                printf("warm[%d] %.3f ms\n", r,
                       (double)(rd_clock() - start) / 1000.0);
                for (i = 0; i < member_cnt; i++)
                        ut_set_owned(&members[i]);
        }

        for (i = 0; i < member_cnt; i++)
                rd_kafka_group_member_clear(&members[i]);
        free(members);
        rd_kafka_metadata_destroy(metadata);
        rd_kafka_destroy(rk);
        return 0;
}

@kphelps
Kyle Phelps (kphelps) requested a review from a team as a code owner July 17, 2026 19:14
@confluent-cla-assistant

confluent-cla-assistant Bot commented Jul 17, 2026

Copy link
Copy Markdown

🎉 All Contributor License Agreements have been signed. Ready to merge.
✅ kphelps
Please push an empty commit if you would like to re-run the checks to verify CLA status for all contributors.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant