DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
| Table of Contents |
|---|
Status
Current state: Under Discussion
...
JIRA:
| Jira | ||||||
|---|---|---|---|---|---|---|
|
Motivation
Kafka allows users to configure how tool scripts behave through multiple methods:
...
Furthermore, the validation of required arguments, such as --bootstrap-server and --bootstrap-controller, should consider multiple sources. For most existing tool scripts, these arguments are only validated when provided via the command line. However, they could also be specified through properties passed as key=value pairs via the command line or configuration files.
Public Interfaces
1) The precedence order for loading properties into Kafka configurations is as follows:
...
- Name: modern
- Description: This configuration determines whether to enable the property loading precedence proposed in this KIP. It will be deprecated in Kafka 5.0, after which the default behavior will always follow the approach proposed by this KIP. It will then be removed in Kafka 6.0.
- Type: boolean
- Default: false
Validating Required Arguments from Multiple Sources
Required arguments (e.g., --bootstrap-server and --bootstrap-controller for now) should be validated after considering values from the following sources:
...
An error should only be raised if the arguments are missing or invalid after checking all three sources.
Proposed Changes
To ensure that all tool scripts correctly honor the proposed property loading precedence, introduce the helper methods in CommandLineUtils.java and enforce their use across all affected tools. This prevents each script from implementing its own logic and ensures consistent behavior.
Since the codebase currently uses three different argument parsers (joptsimple, argparse4j, and AbstractConnectCli), the helper methods are designed to be simple and generic, without depending on any specific external library.
Helper methods
The following code demonstrates the key helper methods:
| Code Block | ||||
|---|---|---|---|---|
| ||||
/**
* Merge multiple configuration sources according to priority, from highest to lowest:
* 1) Command line arguments
* 2) Properties passed as key=value pairs via command line
* 3) Configuration files
* 4) Default values set by tool scripts
*/
public static Map<String, Object> mergePropertiesWithPrecedence(
Map<String, Object> commandLineMap,
Map<String, Object> commandLineKeyValMap,
Map<String, Object> configMap,
Map<String, Object> toolDefaultMap
) {
Map<String, Object> map = new HashMap<>();
// Default values set by tool scripts and default values of command line arguments
if (toolDefaultMap != null) {
map.putAll(toolDefaultMap);
}
// Configuration file
if (configMap != null) {
map.putAll(configMap);
}
// Properties passed as key=value pairs via command line
if (commandLineKeyValMap != null) {
map.putAll(commandLineKeyValMap);
}
// Command line arguments
if (commandLineMap != null) {
map.putAll(commandLineMap);
}
return map;
} |
Example usage
| Code Block | ||||
|---|---|---|---|---|
| ||||
Map<String, Object> readerProps() throws IOException {
Map<String, Object> commandLineMap = new HashMap<>();
commandLineMap.put("topic", options.valueOf(topicOpt));
Map<String, Object> configMap = new HashMap<>();
if (options.has(readerConfigOpt)) {
configMap.putAll(propsToStringMap(loadProps(options.valueOf(readerConfigOpt))));
}
Map<String, Object> commandLineKeyValMap = new HashMap<>();
if (options.has(readerPropertyOpt)) {
commandLineKeyValMap.putAll(propsToStringMap(parseKeyValueArgs(options.valuesOf(readerPropertyOpt))));
}
return mergePropertiesWithPrecedence(commandLineMap, commandLineKeyValMap, configMap, null);
}
Map<String, Object> producerProps() throws IOException {
// Prepare the map from command line arguments
Map<String, Object> commandLineMap = new HashMap<>();
commandLineMap.put(BOOTSTRAP_SERVERS_CONFIG, options.valueOf(bootstrapServerOpt));
commandLineMap.put(COMPRESSION_TYPE_CONFIG, compressionCodec());
if (options.has(sendTimeoutOpt)) commandLineMap.put(LINGER_MS_CONFIG, options.valueOf(sendTimeoutOpt).toString());
if (options.has(requestRequiredAcksOpt)) commandLineMap.put(ACKS_CONFIG, options.valueOf(requestRequiredAcksOpt).toString());
if (options.has(requestTimeoutMsOpt)) commandLineMap.put(REQUEST_TIMEOUT_MS_CONFIG, options.valueOf(requestTimeoutMsOpt).toString());
if (options.has(messageSendMaxRetriesOpt)) commandLineMap.put(RETRIES_CONFIG, options.valueOf(messageSendMaxRetriesOpt).toString());
if (options.has(retryBackoffMsOpt)) commandLineMap.put(RETRY_BACKOFF_MS_CONFIG, options.valueOf(retryBackoffMsOpt).toString());
if (options.has(socketBufferSizeOpt)) commandLineMap.put(SEND_BUFFER_CONFIG, options.valueOf(socketBufferSizeOpt).toString());
if (options.has(maxMemoryBytesOpt)) commandLineMap.put(BUFFER_MEMORY_CONFIG, options.valueOf(maxMemoryBytesOpt).toString());
if (options.has(batchSizeOpt)) commandLineMap.put(BATCH_SIZE_CONFIG, options.valueOf(batchSizeOpt).toString());
if (options.has(maxPartitionMemoryBytesOpt)) commandLineMap.put(BATCH_SIZE_CONFIG, options.valueOf(maxPartitionMemoryBytesOpt).toString());
if (options.has(metadataExpiryMsOpt)) commandLineMap.put(METADATA_MAX_AGE_CONFIG, options.valueOf(metadataExpiryMsOpt).toString());
if (options.has(maxBlockMsOpt)) commandLineMap.put(MAX_BLOCK_MS_CONFIG, options.valueOf(maxBlockMsOpt).toString());
// Properties passed as key=value pairs via command line
Map<String, Object> commandLineKeyValMap = new HashMap<>();
if (options.has(commandPropertyOpt)) {
commandLineKeyValMap.putAll(Utils.propsToStringMap(
parseKeyValueArgs(options.valuesOf(commandPropertyOpt))
));
}
// Configuration files
Map<String, Object> configMap = new HashMap<>();
if (options.has(commandConfigOpt)) {
configMap.putAll(Utils.propsToStringMap(
Utils.loadProps(options.valueOf(commandConfigOpt))
));
}
// Default values set by tool scripts and default values of command line arguments
Map<String, Object> toolDefaultMap = new HashMap<>();
toolDefaultMap.put(CLIENT_ID_CONFIG, "console-producer");
toolDefaultMap.put(KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
toolDefaultMap.put(VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
// Since all the options below have default values, we don't need to check for null
toolDefaultMap.put(LINGER_MS_CONFIG, options.valuesOf(sendTimeoutOpt).toString());
toolDefaultMap.put(ACKS_CONFIG, options.valuesOf(requestRequiredAcksOpt).toString());
toolDefaultMap.put(REQUEST_TIMEOUT_MS_CONFIG, options.valuesOf(requestTimeoutMsOpt).toString());
toolDefaultMap.put(RETRIES_CONFIG, options.valuesOf(messageSendMaxRetriesOpt).toString());
toolDefaultMap.put(RETRY_BACKOFF_MS_CONFIG, options.valuesOf(retryBackoffMsOpt).toString());
toolDefaultMap.put(SEND_BUFFER_CONFIG, options.valuesOf(socketBufferSizeOpt).toString());
toolDefaultMap.put(BUFFER_MEMORY_CONFIG, options.valuesOf(maxMemoryBytesOpt).toString());
toolDefaultMap.put(BATCH_SIZE_CONFIG, options.valuesOf(batchSizeOpt).toString());
toolDefaultMap.put(BATCH_SIZE_CONFIG, options.valuesOf(maxPartitionMemoryBytesOpt).toString());
toolDefaultMap.put(METADATA_MAX_AGE_CONFIG, options.valuesOf(metadataExpiryMsOpt).toString());
toolDefaultMap.put(MAX_BLOCK_MS_CONFIG, options.valuesOf(maxBlockMsOpt).toString());
return mergePropertiesWithPrecedence(commandLineMap, commandLineKeyValMap, configMap, toolDefaultMap);
} |
Required Updates to Existing Tools
The following table lists the tools that need to be updated based on the proposed changes.
...
*New "modern" Option for Deprecation - Indicates that the tool script currently has properties that do not follow the proposed precedence.
Compatibility, Deprecation, and Migration Plan
Impact: Existing users who have configured their settings based on the current tool implementations, rather than the proposed precedence, may be affected.
Deprecation: Adjusting property loading precedence might break existing users. For this reason, we have an option added to each affected tool script to enable or disable this change. The plan is to always use the proposed precedence in the next major release, Kafka 5.0. Concurrently, the option will be deprecated in Kafka 5.0 and subsequently removed in Kafka 6.0.
Test Plan
Unit tests will be added to ensure that the properties are loaded according to the precedence order as proposed.
Rejected Alternatives
1. Deprecate the "modern" option during the proposed rollout and remove it in the next major release, Kafka 5.0
...