Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -628,7 +628,7 @@ private boolean assignTasksToClients(final Cluster fullMetadata,

log.info("{} client nodes and {} consumers participating in this rebalance: \n{}.",
clientStates.size(),
clientStates.values().stream().map(ClientState::capacity).reduce(Integer::sum),
clientStates.values().stream().map(ClientState::capacity).reduce(Integer::sum).orElse(0),

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor fix: It's an Optional and the log currently say:

5 client nodes and Optional[10] consumers participating in this rebalance

clientStates.entrySet().stream()
.sorted(comparingByKey())
.map(entry -> entry.getKey() + ": " + entry.getValue().consumers())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,73 +68,73 @@
@Category({IntegrationTest.class})
@RunWith(MockitoJUnitRunner.StrictStubs.class)
public class StreamsAssignmentScaleTest {
final static long MAX_ASSIGNMENT_DURATION = 60 * 1000L; //each individual assignment should complete within 20s
final static long MAX_ASSIGNMENT_DURATION = 120 * 1000L; // we should stay below `max.poll.interval.ms`
final static String APPLICATION_ID = "streams-assignment-scale-test";

private final Logger log = LoggerFactory.getLogger(StreamsAssignmentScaleTest.class);

/* HighAvailabilityTaskAssignor tests */

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testHighAvailabilityTaskAssignorLargePartitionCount() {
completeLargeAssignment(6_000, 2, 1, 1, HighAvailabilityTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testHighAvailabilityTaskAssignorLargeNumConsumers() {
completeLargeAssignment(1_000, 1_000, 1, 1, HighAvailabilityTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testHighAvailabilityTaskAssignorManyStandbys() {
completeLargeAssignment(1_000, 100, 1, 50, HighAvailabilityTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testHighAvailabilityTaskAssignorManyThreadsPerClient() {
completeLargeAssignment(1_000, 10, 1000, 1, HighAvailabilityTaskAssignor.class);
}

/* StickyTaskAssignor tests */

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testStickyTaskAssignorLargePartitionCount() {
completeLargeAssignment(2_000, 2, 1, 1, StickyTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testStickyTaskAssignorLargeNumConsumers() {
completeLargeAssignment(1_000, 1_000, 1, 1, StickyTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testStickyTaskAssignorManyStandbys() {
completeLargeAssignment(1_000, 100, 1, 20, StickyTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testStickyTaskAssignorManyThreadsPerClient() {
completeLargeAssignment(1_000, 10, 1000, 1, StickyTaskAssignor.class);
}

/* FallbackPriorTaskAssignor tests */

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testFallbackPriorTaskAssignorLargePartitionCount() {
completeLargeAssignment(2_000, 2, 1, 1, FallbackPriorTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testFallbackPriorTaskAssignorLargeNumConsumers() {
completeLargeAssignment(1_000, 1_000, 1, 1, FallbackPriorTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testFallbackPriorTaskAssignorManyStandbys() {
completeLargeAssignment(1_000, 100, 1, 20, FallbackPriorTaskAssignor.class);
}

@Test(timeout = 120 * 1000)
@Test(timeout = 300 * 1000)
public void testFallbackPriorTaskAssignorManyThreadsPerClient() {
completeLargeAssignment(1_000, 10, 1000, 1, FallbackPriorTaskAssignor.class);
}
Expand Down