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 @@ -112,8 +112,6 @@
import org.apache.kafka.common.requests.OffsetFetchResponse;
import org.apache.kafka.common.requests.RenewDelegationTokenRequest;
import org.apache.kafka.common.requests.RenewDelegationTokenResponse;
import org.apache.kafka.common.requests.Resource;
import org.apache.kafka.common.requests.ResourceType;
import org.apache.kafka.common.security.token.delegation.DelegationToken;
import org.apache.kafka.common.security.token.delegation.TokenInformation;
import org.apache.kafka.common.utils.AppInfoParser;
Expand Down Expand Up @@ -1683,19 +1681,19 @@ public DescribeConfigsResult describeConfigs(Collection<ConfigResource> configRe

// The BROKER resources which we want to describe. We must make a separate DescribeConfigs
// request for every BROKER resource we want to describe.
final Collection<Resource> brokerResources = new ArrayList<>();
final Collection<ConfigResource> brokerResources = new ArrayList<>();

// The non-BROKER resources which we want to describe. These resources can be described by a
// single, unified DescribeConfigs request.
final Collection<Resource> unifiedRequestResources = new ArrayList<>(configResources.size());
final Collection<ConfigResource> unifiedRequestResources = new ArrayList<>(configResources.size());

for (ConfigResource resource : configResources) {
if (resource.type() == ConfigResource.Type.BROKER && !resource.isDefault()) {
brokerFutures.put(resource, new KafkaFutureImpl<Config>());
brokerResources.add(configResourceToResource(resource));
brokerFutures.put(resource, new KafkaFutureImpl<>());
brokerResources.add(resource);
} else {
unifiedRequestFutures.put(resource, new KafkaFutureImpl<Config>());
unifiedRequestResources.add(configResourceToResource(resource));
unifiedRequestFutures.put(resource, new KafkaFutureImpl<>());
unifiedRequestResources.add(resource);
}
}

Expand All @@ -1716,7 +1714,7 @@ void handleResponse(AbstractResponse abstractResponse) {
for (Map.Entry<ConfigResource, KafkaFutureImpl<Config>> entry : unifiedRequestFutures.entrySet()) {
ConfigResource configResource = entry.getKey();
KafkaFutureImpl<Config> future = entry.getValue();
DescribeConfigsResponse.Config config = response.config(configResourceToResource(configResource));
DescribeConfigsResponse.Config config = response.config(configResource);
if (config == null) {
future.completeExceptionally(new UnknownServerException(
"Malformed broker response: missing config for " + configResource));
Expand Down Expand Up @@ -1746,7 +1744,7 @@ void handleFailure(Throwable throwable) {

for (Map.Entry<ConfigResource, KafkaFutureImpl<Config>> entry : brokerFutures.entrySet()) {
final KafkaFutureImpl<Config> brokerFuture = entry.getValue();
final Resource resource = configResourceToResource(entry.getKey());
final ConfigResource resource = entry.getKey();
final int nodeId = Integer.parseInt(resource.name());
runnable.call(new Call("describeBrokerConfigs", calcDeadlineMs(now, options.timeoutMs()),
new ConstantNodeIdProvider(nodeId)) {
Expand Down Expand Up @@ -1792,21 +1790,6 @@ void handleFailure(Throwable throwable) {
return new DescribeConfigsResult(allFutures);
}

private Resource configResourceToResource(ConfigResource configResource) {
ResourceType resourceType;
switch (configResource.type()) {
case TOPIC:
resourceType = ResourceType.TOPIC;
break;
case BROKER:
resourceType = ResourceType.BROKER;
break;
default:
throw new IllegalArgumentException("Unexpected resource type " + configResource.type());
}
return new Resource(resourceType, configResource.name());
}

private List<ConfigEntry.ConfigSynonym> configSynonyms(DescribeConfigsResponse.ConfigEntry configEntry) {
List<ConfigEntry.ConfigSynonym> synonyms = new ArrayList<>(configEntry.synonyms().size());
for (DescribeConfigsResponse.ConfigSynonym synonym : configEntry.synonyms()) {
Expand Down Expand Up @@ -1856,21 +1839,21 @@ public AlterConfigsResult alterConfigs(Map<ConfigResource, Config> configs, fina
}
if (!unifiedRequestResources.isEmpty())
allFutures.putAll(alterConfigs(configs, options, unifiedRequestResources, new LeastLoadedNodeProvider()));
return new AlterConfigsResult(new HashMap<ConfigResource, KafkaFuture<Void>>(allFutures));
return new AlterConfigsResult(new HashMap<>(allFutures));
}

private Map<ConfigResource, KafkaFutureImpl<Void>> alterConfigs(Map<ConfigResource, Config> configs,
final AlterConfigsOptions options,
Collection<ConfigResource> resources,
NodeProvider nodeProvider) {
final Map<ConfigResource, KafkaFutureImpl<Void>> futures = new HashMap<>();
final Map<Resource, AlterConfigsRequest.Config> requestMap = new HashMap<>(resources.size());
final Map<ConfigResource, AlterConfigsRequest.Config> requestMap = new HashMap<>(resources.size());
for (ConfigResource resource : resources) {
List<AlterConfigsRequest.ConfigEntry> configEntries = new ArrayList<>();
for (ConfigEntry configEntry: configs.get(resource).entries())
configEntries.add(new AlterConfigsRequest.ConfigEntry(configEntry.name(), configEntry.value()));
requestMap.put(configResourceToResource(resource), new AlterConfigsRequest.Config(configEntries));
futures.put(resource, new KafkaFutureImpl<Void>());
requestMap.put(resource, new AlterConfigsRequest.Config(configEntries));
futures.put(resource, new KafkaFutureImpl<>());
}

final long now = time.milliseconds();
Expand All @@ -1886,7 +1869,7 @@ public void handleResponse(AbstractResponse abstractResponse) {
AlterConfigsResponse response = (AlterConfigsResponse) abstractResponse;
for (Map.Entry<ConfigResource, KafkaFutureImpl<Void>> entry : futures.entrySet()) {
KafkaFutureImpl<Void> future = entry.getValue();
ApiException exception = response.errors().get(configResourceToResource(entry.getKey())).exception();
ApiException exception = response.errors().get(entry.getKey()).exception();
if (exception != null) {
future.completeExceptionally(exception);
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,12 @@

package org.apache.kafka.common.config;

import java.util.Arrays;
import java.util.Collections;
import java.util.Map;
import java.util.Objects;
import java.util.function.Function;
import java.util.stream.Collectors;

/**
* A class representing resources that have configs.
Expand All @@ -28,7 +33,25 @@ public final class ConfigResource {
* Type of resource.
*/
public enum Type {
BROKER, TOPIC, UNKNOWN;
BROKER((byte) 3), TOPIC((byte) 2), UNKNOWN((byte) 0);

private static final Map<Byte, Type> TYPES = Collections.unmodifiableMap(
Arrays.stream(values()).collect(Collectors.toMap(Type::id, Function.identity()))
);

private final byte id;

Type(final byte id) {
this.id = id;
}

public byte id() {
return id;
}

public static Type forId(final byte id) {
return TYPES.getOrDefault(id, UNKNOWN);
}
}

private final Type type;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,12 @@ public Short getOrElse(Field.Int16 field, short alternative) {
return alternative;
}

public Byte getOrElse(Field.Int8 field, byte alternative) {
if (hasField(field.name))
return getByte(field.name);
return alternative;
}

public Integer getOrElse(Field.Int32 field, int alternative) {
if (hasField(field.name))
return getInt(field.name);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

package org.apache.kafka.common.requests;

import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.types.ArrayOf;
import org.apache.kafka.common.protocol.types.Field;
Expand All @@ -29,6 +30,7 @@
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;

import static org.apache.kafka.common.protocol.types.Type.BOOLEAN;
import static org.apache.kafka.common.protocol.types.Type.INT8;
Expand Down Expand Up @@ -73,7 +75,7 @@ public static class Config {
private final Collection<ConfigEntry> entries;

public Config(Collection<ConfigEntry> entries) {
this.entries = entries;
this.entries = Objects.requireNonNull(entries, "entries");
}

public Collection<ConfigEntry> entries() {
Expand All @@ -86,8 +88,8 @@ public static class ConfigEntry {
private final String value;

public ConfigEntry(String name, String value) {
this.name = name;
this.value = value;
this.name = Objects.requireNonNull(name, "name");
this.value = Objects.requireNonNull(value, "value");
}

public String name() {
Expand All @@ -102,12 +104,12 @@ public String value() {

public static class Builder extends AbstractRequest.Builder {

private final Map<Resource, Config> configs;
private final Map<ConfigResource, Config> configs;
private final boolean validateOnly;

public Builder(Map<Resource, Config> configs, boolean validateOnly) {
public Builder(Map<ConfigResource, Config> configs, boolean validateOnly) {
super(ApiKeys.ALTER_CONFIGS);
this.configs = configs;
this.configs = Objects.requireNonNull(configs, "configs");
this.validateOnly = validateOnly;
}

Expand All @@ -117,12 +119,12 @@ public AlterConfigsRequest build(short version) {
}
}

private final Map<Resource, Config> configs;
private final Map<ConfigResource, Config> configs;
private final boolean validateOnly;

public AlterConfigsRequest(short version, Map<Resource, Config> configs, boolean validateOnly) {
public AlterConfigsRequest(short version, Map<ConfigResource, Config> configs, boolean validateOnly) {
super(version);
this.configs = configs;
this.configs = Objects.requireNonNull(configs, "configs");
this.validateOnly = validateOnly;
}

Expand All @@ -134,9 +136,9 @@ public AlterConfigsRequest(Struct struct, short version) {
for (Object resourcesObj : resourcesArray) {
Struct resourcesStruct = (Struct) resourcesObj;

ResourceType resourceType = ResourceType.forId(resourcesStruct.getByte(RESOURCE_TYPE_KEY_NAME));
ConfigResource.Type resourceType = ConfigResource.Type.forId(resourcesStruct.getByte(RESOURCE_TYPE_KEY_NAME));
String resourceName = resourcesStruct.getString(RESOURCE_NAME_KEY_NAME);
Resource resource = new Resource(resourceType, resourceName);
ConfigResource resource = new ConfigResource(resourceType, resourceName);

Object[] configEntriesArray = resourcesStruct.getArray(CONFIG_ENTRIES_KEY_NAME);
List<ConfigEntry> configEntries = new ArrayList<>(configEntriesArray.length);
Expand All @@ -151,7 +153,7 @@ public AlterConfigsRequest(Struct struct, short version) {
}
}

public Map<Resource, Config> configs() {
public Map<ConfigResource, Config> configs() {
return configs;
}

Expand All @@ -164,10 +166,10 @@ protected Struct toStruct() {
Struct struct = new Struct(ApiKeys.ALTER_CONFIGS.requestSchema(version()));
struct.set(VALIDATE_ONLY_KEY_NAME, validateOnly);
List<Struct> resourceStructs = new ArrayList<>(configs.size());
for (Map.Entry<Resource, Config> entry : configs.entrySet()) {
for (Map.Entry<ConfigResource, Config> entry : configs.entrySet()) {
Struct resourceStruct = struct.instance(RESOURCES_KEY_NAME);

Resource resource = entry.getKey();
ConfigResource resource = entry.getKey();
resourceStruct.set(RESOURCE_TYPE_KEY_NAME, resource.type().id());
resourceStruct.set(RESOURCE_NAME_KEY_NAME, resource.name());

Expand All @@ -194,8 +196,8 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
case 0:
case 1:
ApiError error = ApiError.fromThrowable(e);
Map<Resource, ApiError> errors = new HashMap<>(configs.size());
for (Resource resource : configs.keySet())
Map<ConfigResource, ApiError> errors = new HashMap<>(configs.size());
for (ConfigResource resource : configs.keySet())
errors.put(resource, error);
return new AlterConfigsResponse(throttleTimeMs, errors);
default:
Expand All @@ -207,5 +209,4 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable e) {
public static AlterConfigsRequest parse(ByteBuffer buffer, short version) {
return new AlterConfigsRequest(ApiKeys.ALTER_CONFIGS.parseRequest(version, buffer), version);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,20 @@

package org.apache.kafka.common.requests;

import org.apache.kafka.common.config.ConfigResource;
import org.apache.kafka.common.protocol.ApiKeys;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.protocol.types.ArrayOf;
import org.apache.kafka.common.protocol.types.Field;
import org.apache.kafka.common.protocol.types.Schema;
import org.apache.kafka.common.protocol.Errors;
import org.apache.kafka.common.protocol.types.Struct;

import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;

import static org.apache.kafka.common.protocol.CommonFields.ERROR_CODE;
import static org.apache.kafka.common.protocol.CommonFields.ERROR_MESSAGE;
Expand Down Expand Up @@ -62,12 +64,11 @@ public static Schema[] schemaVersions() {
}

private final int throttleTimeMs;
private final Map<Resource, ApiError> errors;
private final Map<ConfigResource, ApiError> errors;

public AlterConfigsResponse(int throttleTimeMs, Map<Resource, ApiError> errors) {
public AlterConfigsResponse(int throttleTimeMs, Map<ConfigResource, ApiError> errors) {
this.throttleTimeMs = throttleTimeMs;
this.errors = errors;

this.errors = Objects.requireNonNull(errors, "errors");
}

public AlterConfigsResponse(Struct struct) {
Expand All @@ -77,13 +78,13 @@ public AlterConfigsResponse(Struct struct) {
for (Object resourceObj : resourcesArray) {
Struct resourceStruct = (Struct) resourceObj;
ApiError error = new ApiError(resourceStruct);
ResourceType resourceType = ResourceType.forId(resourceStruct.getByte(RESOURCE_TYPE_KEY_NAME));
ConfigResource.Type resourceType = ConfigResource.Type.forId(resourceStruct.getByte(RESOURCE_TYPE_KEY_NAME));
String resourceName = resourceStruct.getString(RESOURCE_NAME_KEY_NAME);
errors.put(new Resource(resourceType, resourceName), error);
errors.put(new ConfigResource(resourceType, resourceName), error);
}
}

public Map<Resource, ApiError> errors() {
public Map<ConfigResource, ApiError> errors() {
return errors;
}

Expand All @@ -102,9 +103,9 @@ protected Struct toStruct(short version) {
Struct struct = new Struct(ApiKeys.ALTER_CONFIGS.responseSchema(version));
struct.set(THROTTLE_TIME_MS, throttleTimeMs);
List<Struct> resourceStructs = new ArrayList<>(errors.size());
for (Map.Entry<Resource, ApiError> entry : errors.entrySet()) {
for (Map.Entry<ConfigResource, ApiError> entry : errors.entrySet()) {
Struct resourceStruct = struct.instance(RESOURCES_KEY_NAME);
Resource resource = entry.getKey();
ConfigResource resource = entry.getKey();
entry.getValue().write(resourceStruct);
resourceStruct.set(RESOURCE_TYPE_KEY_NAME, resource.type().id());
resourceStruct.set(RESOURCE_NAME_KEY_NAME, resource.name());
Expand Down
Loading