Skip to content
Closed
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 @@ -650,7 +650,7 @@ public void transitionToSendingLeaveGroup(boolean dueToExpiredPollTimer) {
* the group, this will be invoked with empty epoch.
*/
void notifyEpochChange(Optional<Integer> epoch) {
stateUpdatesListeners.forEach(stateListener -> stateListener.onMemberEpochUpdated(epoch, memberId));
stateUpdatesListeners.forEach(stateListener -> stateListener.onMemberEpochUpdated(epoch, Optional.of(memberId)));
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -266,7 +266,7 @@ private void process(final ConsumerRebalanceListenerCallbackNeededEvent event) {

private final MemberStateListener memberStateListener = new MemberStateListener() {
@Override
public void onMemberEpochUpdated(Optional<Integer> memberEpoch, String memberId) {
public void onMemberEpochUpdated(Optional<Integer> memberEpoch, Optional<String> memberId) {
updateGroupMetadata(memberEpoch, memberId);
}

Expand Down Expand Up @@ -387,7 +387,9 @@ public AsyncKafkaConsumer(final ConsumerConfig config,
applicationEventProcessorSupplier,
networkClientDelegateSupplier,
requestManagersSupplier);

streamsAssignmentInterface.ifPresent(
sai -> sai.setApplicationEventHandler(applicationEventHandler)
);
this.rebalanceListenerInvoker = new ConsumerRebalanceListenerInvoker(
logContext,
subscriptions,
Expand Down Expand Up @@ -477,8 +479,7 @@ public AsyncKafkaConsumer(final ConsumerConfig config,
Deserializer<V> valueDeserializer,
KafkaClient client,
SubscriptionState subscriptions,
ConsumerMetadata metadata,
Optional<StreamsAssignmentInterface> streamsInstanceMetadata) {
ConsumerMetadata metadata) {
this.log = logContext.logger(getClass());
this.subscriptions = subscriptions;
this.clientId = config.getString(ConsumerConfig.CLIENT_ID_CONFIG);
Expand Down Expand Up @@ -546,8 +547,8 @@ public AsyncKafkaConsumer(final ConsumerConfig config,
clientTelemetryReporter,
metrics,
offsetCommitCallbackInvoker,
memberStateListener,
streamsInstanceMetadata
this::updateGroupMetadata,
Optional.empty()
);
Supplier<ApplicationEventProcessor> applicationEventProcessorSupplier = ApplicationEventProcessor.supplier(
logContext,
Expand Down Expand Up @@ -651,13 +652,13 @@ private ConsumerGroupMetadata initializeConsumerGroupMetadata(final String group
);
}

private void updateGroupMetadata(final Optional<Integer> memberEpoch, final String memberId) {
private void updateGroupMetadata(final Optional<Integer> memberEpoch, final Optional<String> memberId) {
memberEpoch.ifPresent(epoch -> groupMetadata.updateAndGet(
oldGroupMetadataOptional -> oldGroupMetadataOptional.map(
oldGroupMetadata -> new ConsumerGroupMetadata(
oldGroupMetadata.groupId(),
memberEpoch.orElse(oldGroupMetadata.generationId()),
memberId,
memberId.orElse(oldGroupMetadata.memberId()),
oldGroupMetadata.groupInstanceId()
)
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,8 +92,7 @@ public <K, V> ConsumerDelegate<K, V> create(LogContext logContext,
valueDeserializer,
client,
subscriptions,
metadata,
Optional.empty()
metadata
);
else
return new ClassicKafkaConsumer<>(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ public interface MemberStateListener {
* not part of the group anymore.
* @param memberId Current member ID. It won't change until the process is terminated.
*/
void onMemberEpochUpdated(Optional<Integer> memberEpoch, String memberId);
void onMemberEpochUpdated(Optional<Integer> memberEpoch, Optional<String> memberId);

/**
* This callback is invoked when a group member's assigned set of partitions changes. Assignments can change via
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ public class RequestManagers implements Closeable {
public final Optional<ShareHeartbeatRequestManager> shareHeartbeatRequestManager;
public final Optional<ConsumerMembershipManager> consumerMembershipManager;
public final Optional<ShareMembershipManager> shareMembershipManager;
public final Optional<StreamsMembershipManager> streamsMembershipManager;
public final OffsetsRequestManager offsetsRequestManager;
public final TopicMetadataRequestManager topicMetadataRequestManager;
public final FetchRequestManager fetchRequestManager;
Expand All @@ -70,7 +71,8 @@ public RequestManagers(LogContext logContext,
Optional<CommitRequestManager> commitRequestManager,
Optional<ConsumerHeartbeatRequestManager> heartbeatRequestManager,
Optional<ConsumerMembershipManager> membershipManager,
Optional<StreamsGroupHeartbeatRequestManager> streamsGroupHeartbeatRequestManager) {
Optional<StreamsGroupHeartbeatRequestManager> streamsGroupHeartbeatRequestManager,
Optional<StreamsMembershipManager> streamsMembershipManager) {
this.log = logContext.logger(RequestManagers.class);
this.offsetsRequestManager = requireNonNull(offsetsRequestManager, "OffsetsRequestManager cannot be null");
this.coordinatorRequestManager = coordinatorRequestManager;
Expand All @@ -83,6 +85,7 @@ public RequestManagers(LogContext logContext,
this.consumerMembershipManager = membershipManager;
this.shareMembershipManager = Optional.empty();
this.streamsGroupHeartbeatRequestManager = streamsGroupHeartbeatRequestManager;
this.streamsMembershipManager = streamsMembershipManager;

List<Optional<? extends RequestManager>> list = new ArrayList<>();
list.add(coordinatorRequestManager);
Expand All @@ -93,6 +96,7 @@ public RequestManagers(LogContext logContext,
list.add(Optional.of(topicMetadataRequestManager));
list.add(Optional.of(fetchRequestManager));
list.add(streamsGroupHeartbeatRequestManager);
list.add(streamsMembershipManager);
entries = Collections.unmodifiableList(list);
}

Expand All @@ -110,6 +114,7 @@ public RequestManagers(LogContext logContext,
this.consumerMembershipManager = Optional.empty();
this.shareMembershipManager = shareMembershipManager;
this.streamsGroupHeartbeatRequestManager = Optional.empty();
this.streamsMembershipManager = Optional.empty();
this.offsetsRequestManager = null;
this.topicMetadataRequestManager = null;
this.fetchRequestManager = null;
Expand Down Expand Up @@ -194,6 +199,7 @@ protected RequestManagers create() {
CoordinatorRequestManager coordinator = null;
CommitRequestManager commitRequestManager = null;
StreamsGroupHeartbeatRequestManager streamsGroupHeartbeatRequestManager = null;
StreamsMembershipManager streamsMembershipManager = null;

if (groupRebalanceConfig != null && groupRebalanceConfig.groupId != null) {
Optional<String> serverAssignor = Optional.ofNullable(config.getString(ConsumerConfig.GROUP_REMOTE_ASSIGNOR_CONFIG));
Expand Down Expand Up @@ -240,17 +246,18 @@ protected RequestManagers create() {

if (streamsInstanceMetadata.isPresent()) {
streamsGroupHeartbeatRequestManager = new StreamsGroupHeartbeatRequestManager(
logContext,
time,
config,
coordinator,
membershipManager,
backgroundEventHandler,
metrics,
streamsInstanceMetadata.get(),
metadata
logContext,
time,
config,
coordinator,
streamsMembershipManager,
backgroundEventHandler,
metrics,
streamsInstanceMetadata.get()
);
} else {
membershipManager.registerStateListener(commitRequestManager);
membershipManager.registerStateListener(applicationThreadMemberStateListener);
heartbeatRequestManager = new ConsumerHeartbeatRequestManager(
logContext,
time,
Expand Down Expand Up @@ -284,7 +291,8 @@ protected RequestManagers create() {
Optional.ofNullable(commitRequestManager),
Optional.ofNullable(heartbeatRequestManager),
Optional.ofNullable(membershipManager),
Optional.ofNullable(streamsGroupHeartbeatRequestManager)
Optional.ofNullable(streamsGroupHeartbeatRequestManager),
Optional.ofNullable(streamsMembershipManager)
);
}
};
Expand Down
Loading