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 @@ -872,6 +872,9 @@ public static class InternalConfig {
public static final String ASSIGNMENT_ERROR_CODE = "__assignment.error.code__";
public static final String NEXT_SCHEDULED_REBALANCE_MS = "__next.probing.rebalance.ms__";
public static final String TIME = "__time__";

// This is settable in the main Streams config, but it's a private API for testing
public static final String ASSIGNMENT_LISTENER = "__asignment.listener__";
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.kafka.streams.processor.internals.assignment.AssignmentInfo;
import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration;
import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration.AssignmentConfigs;
import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration.AssignmentListener;
import org.apache.kafka.streams.processor.internals.assignment.AssignorError;
import org.apache.kafka.streams.processor.internals.assignment.ClientState;
import org.apache.kafka.streams.processor.internals.assignment.CopartitionedTopicsEnforcer;
Expand Down Expand Up @@ -171,6 +172,7 @@ public String toString() {
private InternalTopicManager internalTopicManager;
private CopartitionedTopicsEnforcer copartitionedTopicsEnforcer;
private RebalanceProtocol rebalanceProtocol;
private AssignmentListener assignmentListener;

private Supplier<TaskAssignor> taskAssignorSupplier;

Expand All @@ -189,20 +191,21 @@ public void configure(final Map<String, ?> configs) {
log = new LogContext(logPrefix).logger(getClass());
usedSubscriptionMetadataVersion = assignorConfiguration
.configuredMetadataVersion(usedSubscriptionMetadataVersion);
taskManager = assignorConfiguration.getTaskManager();
streamsMetadataState = assignorConfiguration.getStreamsMetadataState();
assignmentErrorCode = assignorConfiguration.getAssignmentErrorCode(configs);
nextScheduledRebalanceMs = assignorConfiguration.getNextScheduledRebalanceMs(configs);
time = assignorConfiguration.getTime(configs);
assignmentConfigs = assignorConfiguration.getAssignmentConfigs();
partitionGrouper = assignorConfiguration.getPartitionGrouper();
userEndPoint = assignorConfiguration.getUserEndPoint();
adminClient = assignorConfiguration.getAdminClient();
adminClientTimeout = assignorConfiguration.getAdminClientTimeout();
internalTopicManager = assignorConfiguration.getInternalTopicManager();
copartitionedTopicsEnforcer = assignorConfiguration.getCopartitionedTopicsEnforcer();
taskManager = assignorConfiguration.taskManager();
streamsMetadataState = assignorConfiguration.streamsMetadataState();
assignmentErrorCode = assignorConfiguration.assignmentErrorCode();
nextScheduledRebalanceMs = assignorConfiguration.nextScheduledRebalanceMs();
time = assignorConfiguration.time();
assignmentConfigs = assignorConfiguration.assignmentConfigs();
partitionGrouper = assignorConfiguration.partitionGrouper();
userEndPoint = assignorConfiguration.userEndPoint();
adminClient = assignorConfiguration.adminClient();
adminClientTimeout = assignorConfiguration.adminClientTimeout();
internalTopicManager = assignorConfiguration.internalTopicManager();
copartitionedTopicsEnforcer = assignorConfiguration.copartitionedTopicsEnforcer();
rebalanceProtocol = assignorConfiguration.rebalanceProtocol();
taskAssignorSupplier = assignorConfiguration::getTaskAssignor;
taskAssignorSupplier = assignorConfiguration::taskAssignor;
assignmentListener = assignorConfiguration.assignmentListener();
}

@Override
Expand Down Expand Up @@ -913,8 +916,10 @@ private Map<String, Assignment> computeNewAssignment(final Map<UUID, ClientMetad
}

if (rebalanceRequired) {
assignmentListener.onAssignmentComplete(false);
log.info("Finished unstable assignment of tasks, a followup rebalance will be scheduled.");
} else {
assignmentListener.onAssignmentComplete(true);
log.info("Finished stable assignment of tasks, no followup rebalances required.");
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,59 +47,21 @@ public final class AssignorConfiguration {

private final String logPrefix;
private final Logger log;
private final AssignmentConfigs assignmentConfigs;
@SuppressWarnings("deprecation")
private final org.apache.kafka.streams.processor.PartitionGrouper partitionGrouper;
private final String userEndPoint;
private final TaskManager taskManager;
private final StreamsMetadataState streamsMetadataState;

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.

Sorry for all the lines of changes here, but the inconsistent method signatures and object construction was making this class hard to follow. All I'm doing here is making the methods conform to the same style, and moving the construction of any config that isn't needed elsewhere to its getter method. This allows us to remove a lot of these class variables

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Do we know if the corresponding "getters" are called often or not? I guess the main idea behind having those member variables was to "parse" the config once as it's immutable anyway instead of each time a getter is called?

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.

They are not. We construct the AssignorConfiguration object once in StreamsPartitionAssignor#configure and call the getters immediately afterwards, and only there

private final Admin adminClient;
private final int adminClientTimeout;
private final InternalTopicManager internalTopicManager;
private final CopartitionedTopicsEnforcer copartitionedTopicsEnforcer;

private final StreamsConfig streamsConfig;
private final Map<String, ?> internalConfigs;

@SuppressWarnings("deprecation")
public AssignorConfiguration(final Map<String, ?> configs) {
streamsConfig = new QuietStreamsConfig(configs);
internalConfigs = configs;

// Setting the logger with the passed in client thread name
logPrefix = String.format("stream-thread [%s] ", streamsConfig.getString(CommonClientConfigs.CLIENT_ID_CONFIG));
final LogContext logContext = new LogContext(logPrefix);
log = logContext.logger(getClass());

assignmentConfigs = new AssignmentConfigs(streamsConfig);

partitionGrouper = streamsConfig.getConfiguredInstance(
StreamsConfig.PARTITION_GROUPER_CLASS_CONFIG,
org.apache.kafka.streams.processor.PartitionGrouper.class
);

final String configuredUserEndpoint = streamsConfig.getString(StreamsConfig.APPLICATION_SERVER_CONFIG);
if (configuredUserEndpoint != null && !configuredUserEndpoint.isEmpty()) {
try {
final String host = getHost(configuredUserEndpoint);
final Integer port = getPort(configuredUserEndpoint);

if (host == null || port == null) {
throw new ConfigException(
String.format(
"%s Config %s isn't in the correct format. Expected a host:port pair but received %s",
logPrefix, StreamsConfig.APPLICATION_SERVER_CONFIG, configuredUserEndpoint
)
);
}
} catch (final NumberFormatException nfe) {
throw new ConfigException(
String.format("%s Invalid port supplied in %s for config %s: %s",
logPrefix, configuredUserEndpoint, StreamsConfig.APPLICATION_SERVER_CONFIG, nfe)
);
}
userEndPoint = configuredUserEndpoint;
} else {
userEndPoint = null;
}

{
final Object o = configs.get(StreamsConfig.InternalConfig.TASK_MANAGER_FOR_PARTITION_ASSIGNOR);
if (o == null) {
Expand All @@ -119,25 +81,6 @@ public AssignorConfiguration(final Map<String, ?> configs) {
taskManager = (TaskManager) o;
}

{
final Object o = configs.get(StreamsConfig.InternalConfig.STREAMS_METADATA_STATE_FOR_PARTITION_ASSIGNOR);
if (o == null) {
final KafkaException fatalException = new KafkaException("StreamsMetadataState is not specified");
log.error(fatalException.getMessage(), fatalException);
throw fatalException;
}

if (!(o instanceof StreamsMetadataState)) {
final KafkaException fatalException = new KafkaException(
String.format("%s is not an instance of %s", o.getClass().getName(), StreamsMetadataState.class.getName())
);
log.error(fatalException.getMessage(), fatalException);
throw fatalException;
}

streamsMetadataState = (StreamsMetadataState) o;
}

{
final Object o = configs.get(StreamsConfig.InternalConfig.STREAMS_ADMIN_CLIENT);
if (o == null) {
Expand All @@ -155,13 +98,8 @@ public AssignorConfiguration(final Map<String, ?> configs) {
}

adminClient = (Admin) o;
internalTopicManager = new InternalTopicManager(adminClient, streamsConfig);
}

adminClientTimeout = streamsConfig.getInt(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG);

copartitionedTopicsEnforcer = new CopartitionedTopicsEnforcer(logPrefix);

{
final String o = (String) configs.get(INTERNAL_TASK_ASSIGNOR_CLASS);
if (o == null) {
Expand All @@ -172,8 +110,8 @@ public AssignorConfiguration(final Map<String, ?> configs) {
}
}

public AtomicInteger getAssignmentErrorCode(final Map<String, ?> configs) {
final Object ai = configs.get(StreamsConfig.InternalConfig.ASSIGNMENT_ERROR_CODE);
public AtomicInteger assignmentErrorCode() {
final Object ai = internalConfigs.get(StreamsConfig.InternalConfig.ASSIGNMENT_ERROR_CODE);
if (ai == null) {
final KafkaException fatalException = new KafkaException("assignmentErrorCode is not specified");
log.error(fatalException.getMessage(), fatalException);
Expand All @@ -190,8 +128,8 @@ public AtomicInteger getAssignmentErrorCode(final Map<String, ?> configs) {
return (AtomicInteger) ai;
}

public AtomicLong getNextScheduledRebalanceMs(final Map<String, ?> configs) {
final Object al = configs.get(InternalConfig.NEXT_SCHEDULED_REBALANCE_MS);
public AtomicLong nextScheduledRebalanceMs() {
final Object al = internalConfigs.get(InternalConfig.NEXT_SCHEDULED_REBALANCE_MS);
if (al == null) {
final KafkaException fatalException = new KafkaException("nextProbingRebalanceMs is not specified");
log.error(fatalException.getMessage(), fatalException);
Expand All @@ -209,8 +147,8 @@ public AtomicLong getNextScheduledRebalanceMs(final Map<String, ?> configs) {
return (AtomicLong) al;
}

public Time getTime(final Map<String, ?> configs) {
final Object t = configs.get(InternalConfig.TIME);
public Time time() {
final Object t = internalConfigs.get(InternalConfig.TIME);
if (t == null) {
final KafkaException fatalException = new KafkaException("time is not specified");
log.error(fatalException.getMessage(), fatalException);
Expand All @@ -228,12 +166,27 @@ public Time getTime(final Map<String, ?> configs) {
return (Time) t;
}

public TaskManager getTaskManager() {
public TaskManager taskManager() {
return taskManager;
}

public StreamsMetadataState getStreamsMetadataState() {
return streamsMetadataState;
public StreamsMetadataState streamsMetadataState() {
final Object o = internalConfigs.get(StreamsConfig.InternalConfig.STREAMS_METADATA_STATE_FOR_PARTITION_ASSIGNOR);
if (o == null) {
final KafkaException fatalException = new KafkaException("StreamsMetadataState is not specified");
log.error(fatalException.getMessage(), fatalException);
throw fatalException;
}

if (!(o instanceof StreamsMetadataState)) {
final KafkaException fatalException = new KafkaException(
String.format("%s is not an instance of %s", o.getClass().getName(), StreamsMetadataState.class.getName())
);
log.error(fatalException.getMessage(), fatalException);
throw fatalException;
}

return (StreamsMetadataState) o;
}

public RebalanceProtocol rebalanceProtocol() {
Expand Down Expand Up @@ -301,35 +254,61 @@ public int configuredMetadataVersion(final int priorVersion) {
}

@SuppressWarnings("deprecation")
public org.apache.kafka.streams.processor.PartitionGrouper getPartitionGrouper() {
return partitionGrouper;
public org.apache.kafka.streams.processor.PartitionGrouper partitionGrouper() {
return streamsConfig.getConfiguredInstance(
StreamsConfig.PARTITION_GROUPER_CLASS_CONFIG,
org.apache.kafka.streams.processor.PartitionGrouper.class
);
}

public String getUserEndPoint() {
return userEndPoint;
public String userEndPoint() {
final String configuredUserEndpoint = streamsConfig.getString(StreamsConfig.APPLICATION_SERVER_CONFIG);
if (configuredUserEndpoint != null && !configuredUserEndpoint.isEmpty()) {
try {
final String host = getHost(configuredUserEndpoint);
final Integer port = getPort(configuredUserEndpoint);

if (host == null || port == null) {
throw new ConfigException(
String.format(
"%s Config %s isn't in the correct format. Expected a host:port pair but received %s",
logPrefix, StreamsConfig.APPLICATION_SERVER_CONFIG, configuredUserEndpoint
)
);
}
} catch (final NumberFormatException nfe) {
throw new ConfigException(
String.format("%s Invalid port supplied in %s for config %s: %s",
logPrefix, configuredUserEndpoint, StreamsConfig.APPLICATION_SERVER_CONFIG, nfe)
);
}
return configuredUserEndpoint;
} else {
return null;
}
}

public Admin getAdminClient() {
public Admin adminClient() {
return adminClient;
}

public int getAdminClientTimeout() {
return adminClientTimeout;
public int adminClientTimeout() {
return streamsConfig.getInt(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG);
}

public InternalTopicManager getInternalTopicManager() {
return internalTopicManager;
public InternalTopicManager internalTopicManager() {
return new InternalTopicManager(adminClient, streamsConfig);
}

public CopartitionedTopicsEnforcer getCopartitionedTopicsEnforcer() {
return copartitionedTopicsEnforcer;
public CopartitionedTopicsEnforcer copartitionedTopicsEnforcer() {
return new CopartitionedTopicsEnforcer(logPrefix);
}

public AssignmentConfigs getAssignmentConfigs() {
return assignmentConfigs;
public AssignmentConfigs assignmentConfigs() {
return new AssignmentConfigs(streamsConfig);
}

public TaskAssignor getTaskAssignor() {
public TaskAssignor taskAssignor() {
try {
return Utils.newInstance(taskAssignorClass, TaskAssignor.class);
} catch (final ClassNotFoundException e) {
Expand All @@ -340,6 +319,27 @@ public TaskAssignor getTaskAssignor() {
}
}

public AssignmentListener assignmentListener() {

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.

This is the only non-cosmetic change in this class, along with the interface added below

final Object o = internalConfigs.get(InternalConfig.ASSIGNMENT_LISTENER);
if (o == null) {
return stable -> { };
}

if (!(o instanceof AssignmentListener)) {
final KafkaException fatalException = new KafkaException(
String.format("%s is not an instance of %s", o.getClass().getName(), AssignmentListener.class.getName())
);
log.error(fatalException.getMessage(), fatalException);
throw fatalException;
}

return (AssignmentListener) o;
}

public interface AssignmentListener {
void onAssignmentComplete(final boolean stable);
}

public static class AssignmentConfigs {
public final long acceptableRecoveryLag;
public final int maxWarmupReplicas;
Expand Down
Loading