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 @@ -66,6 +66,7 @@
import org.slf4j.Logger;

import java.nio.ByteBuffer;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
Expand Down Expand Up @@ -521,14 +522,18 @@ public boolean poll(Timer timer, boolean waitForJoinGroup) {
}
}
} else {
// For manually assigned partitions, if coordinator is unknown, make sure we lookup one and await metadata.
// For manually assigned partitions, we do not try to pro-actively lookup coordinator;
// instead we only try to refresh metadata when necessary.
// If connections to all nodes fail, wakeups triggered while attempting to send fetch
// requests result in polls returning immediately, causing a tight loop of polls. Without
// the wakeup, poll() with no channels would block for the timeout, delaying re-connection.
// awaitMetadataUpdate() in ensureCoordinatorReady initiates new connections with configured backoff and avoids the busy loop.
if (coordinatorUnknownAndUnready(timer)) {
return false;
if (metadata.updateRequested() && !client.hasReadyNodes(timer.currentTimeMs())) {
client.awaitMetadataUpdate(timer);
}

// if there is pending coordinator requests, ensure they have a chance to be transmitted.
client.pollNoWakeup();
}

maybeAutoCommitOffsetsAsync(timer.currentTimeMs());
Expand Down Expand Up @@ -946,7 +951,17 @@ void invokeCompletedOffsetCommitCallbacks() {
public void commitOffsetsAsync(final Map<TopicPartition, OffsetAndMetadata> offsets, final OffsetCommitCallback callback) {
invokeCompletedOffsetCommitCallbacks();

if (!coordinatorUnknown()) {
if (!coordinatorUnknownAndUnready(time.timer(Duration.ZERO))) {
// we need to make sure coordinator is ready before committing, since
// this is for async committing we do not try to block, but just try once to
// clear the previous discover-coordinator future, resend, or get responses;
// if the coordinator is not ready yet then we would just proceed and put that into the
// pending requests, and future poll calls would still try to complete them.
//
// the key here though is that we have to try sending the discover-coordinator if
// it's not known or ready, since this is the only place we can send such request
// under manual assignment (there we would not have heartbeat thread trying to auto-rediscover
// the coordinator).
doCommitOffsetsAsync(offsets, callback);
} else {
// we don't know the current coordinator, so try to find it and then send the commit
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -716,7 +716,6 @@ public void testFetchProgressWithMissingPartitionPosition() {
consumer.seekToEnd(singleton(tp0));
consumer.seekToBeginning(singleton(tp1));

client.prepareResponseFrom(FindCoordinatorResponse.prepareResponse(Errors.NONE, groupId, node), node);
client.prepareResponse(body -> {
ListOffsetsRequest request = (ListOffsetsRequest) body;
List<ListOffsetsPartition> partitions = request.topics().stream().flatMap(t -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -486,10 +486,62 @@ public void testCoordinatorNotAvailableWithUserAssignedType() {
coordinator.poll(time.timer(0));
assertTrue(coordinator.coordinatorUnknown());

// should find an available node in next find coordinator request
// should not try to find coordinator since we are in manual assignment
// hence the prepared response should not be returned
client.prepareResponse(groupCoordinatorResponse(node, Errors.NONE));
coordinator.poll(time.timer(Long.MAX_VALUE));
assertTrue(coordinator.coordinatorUnknown());
}

@Test
public void testAutoCommitAsyncWithUserAssignedType() {
try (ConsumerCoordinator coordinator = buildCoordinator(rebalanceConfig, new Metrics(), assignors, true, subscriptions)) {
subscriptions.assignFromUser(Collections.singleton(t1p));
// set timeout to 0 because we expect no requests sent
coordinator.poll(time.timer(0));
assertTrue(coordinator.coordinatorUnknown());
assertFalse(client.hasInFlightRequests());

// elapse auto commit interval and set committable position
time.sleep(autoCommitIntervalMs);
subscriptions.seekUnvalidated(t1p, new SubscriptionState.FetchPosition(100L));

// should try to find coordinator since we are auto committing
coordinator.poll(time.timer(0));
assertTrue(coordinator.coordinatorUnknown());
assertTrue(client.hasInFlightRequests());

client.respond(groupCoordinatorResponse(node, Errors.NONE));
coordinator.poll(time.timer(0));
assertFalse(coordinator.coordinatorUnknown());
// after we've discovered the coordinator we should send
// out the commit request immediately
assertTrue(client.hasInFlightRequests());
}
}

@Test
public void testCommitAsyncWithUserAssignedType() {
subscriptions.assignFromUser(Collections.singleton(t1p));
// set timeout to 0 because we expect no requests sent
coordinator.poll(time.timer(0));
assertTrue(coordinator.coordinatorUnknown());
assertFalse(client.hasInFlightRequests());

// should try to find coordinator since we are commit async
coordinator.commitOffsetsAsync(singletonMap(t1p, new OffsetAndMetadata(100L)), (offsets, exception) -> {
fail("Commit should not get responses, but got offsets:" + offsets + ", and exception:" + exception);
});
coordinator.poll(time.timer(0));
assertTrue(coordinator.coordinatorUnknown());
assertTrue(client.hasInFlightRequests());

client.respond(groupCoordinatorResponse(node, Errors.NONE));
coordinator.poll(time.timer(0));
assertFalse(coordinator.coordinatorUnknown());
// after we've discovered the coordinator we should send
// out the commit request immediately
assertTrue(client.hasInFlightRequests());
}

@Test
Expand Down Expand Up @@ -1901,8 +1953,7 @@ private void testInFlightRequestsFailedAfterCoordinatorMarkedDead(Errors error)

@Test
public void testAutoCommitDynamicAssignment() {
try (ConsumerCoordinator coordinator = buildCoordinator(rebalanceConfig, new Metrics(), assignors, true, subscriptions)
) {
try (ConsumerCoordinator coordinator = buildCoordinator(rebalanceConfig, new Metrics(), assignors, true, subscriptions)) {
subscriptions.subscribe(singleton(topic1), rebalanceListener);
joinAsFollowerAndReceiveAssignment(coordinator, singletonList(t1p));
subscriptions.seek(t1p, 100);
Expand Down