diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index 124b028eb41c3..2ded53387107a 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -336,9 +336,9 @@
+ files="(LogValidator|RemoteLogManagerConfig|RemoteLogManager).java"/>
+ files="(LogValidator|RemoteLogManager|RemoteIndexCache).java"/>
diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
index 1f82499cc8189..fb1b6268c57d4 100644
--- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
+++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java
@@ -1228,7 +1228,7 @@ private void handleResponses(long now, List responses) {
call.fail(now, authException);
} else {
call.fail(now, new DisconnectException(String.format(
- "Cancelled %s request with correlation id %s due to node %s being disconnected",
+ "Cancelled %s request with correlation id %d due to node %s being disconnected",
call.callName, correlationId, response.destination())));
}
} else {
diff --git a/clients/src/main/java/org/apache/kafka/common/record/MemoryRecordsBuilder.java b/clients/src/main/java/org/apache/kafka/common/record/MemoryRecordsBuilder.java
index 3e9360f04ca69..479b306a1d652 100644
--- a/clients/src/main/java/org/apache/kafka/common/record/MemoryRecordsBuilder.java
+++ b/clients/src/main/java/org/apache/kafka/common/record/MemoryRecordsBuilder.java
@@ -438,7 +438,7 @@ private void appendWithOffset(long offset, boolean isControlRecord, long timesta
throw new IllegalArgumentException("Control records can only be appended to control batches");
if (lastOffset != null && offset <= lastOffset)
- throw new IllegalArgumentException(String.format("Illegal offset %s following previous offset %s " +
+ throw new IllegalArgumentException(String.format("Illegal offset %d following previous offset %d " +
"(Offsets must increase monotonically).", offset, lastOffset));
if (timestamp < 0 && timestamp != RecordBatch.NO_TIMESTAMP)
diff --git a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/HttpAccessTokenRetriever.java b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/HttpAccessTokenRetriever.java
index f0362f00f297f..6544005835b07 100644
--- a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/HttpAccessTokenRetriever.java
+++ b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/HttpAccessTokenRetriever.java
@@ -272,7 +272,7 @@ static String handleOutput(final HttpURLConnection con) throws IOException {
errorResponseBody);
if (responseBody == null || responseBody.isEmpty())
- throw new IOException(String.format("The token endpoint response was unexpectedly empty despite response code %s from %s and error message %s",
+ throw new IOException(String.format("The token endpoint response was unexpectedly empty despite response code %d from %s and error message %s",
responseCode, con.getURL(), formatErrorMessage(errorResponseBody)));
return responseBody;
@@ -337,7 +337,7 @@ static String parseAccessToken(String responseBody) throws IOException {
if (snippet.length() > MAX_RESPONSE_BODY_LENGTH) {
int actualLength = responseBody.length();
String s = responseBody.substring(0, MAX_RESPONSE_BODY_LENGTH);
- snippet = String.format("%s (trimmed to first %s characters out of %s total)", s, MAX_RESPONSE_BODY_LENGTH, actualLength);
+ snippet = String.format("%s (trimmed to first %d characters out of %d total)", s, MAX_RESPONSE_BODY_LENGTH, actualLength);
}
throw new IOException(String.format("The token endpoint response did not contain an access_token value. Response: (%s)", snippet));
diff --git a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/RefreshingHttpsJwks.java b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/RefreshingHttpsJwks.java
index 5dc57dead3a93..590c5d8b7499e 100644
--- a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/RefreshingHttpsJwks.java
+++ b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/RefreshingHttpsJwks.java
@@ -348,8 +348,8 @@ public boolean maybeExpediteRefresh(String keyId) {
// 1. Don't try to resolve the key as the large ID will sit in our cache
// 2. Report the issue in the logs but include only the first N characters
int actualLength = keyId.length();
- String s = keyId.substring(0, MISSING_KEY_ID_MAX_KEY_LENGTH);
- String snippet = String.format("%s (trimmed to first %s characters out of %s total)", s, MISSING_KEY_ID_MAX_KEY_LENGTH, actualLength);
+ String trimmedKeyId = keyId.substring(0, MISSING_KEY_ID_MAX_KEY_LENGTH);
+ String snippet = String.format("%s (trimmed to first %d characters out of %d total)", trimmedKeyId, MISSING_KEY_ID_MAX_KEY_LENGTH, actualLength);
log.warn("Key ID {} was too long to cache", snippet);
return false;
} else {
diff --git a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/SerializedJwt.java b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/SerializedJwt.java
index 6456e8b06c3f6..82c63111d1f2c 100644
--- a/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/SerializedJwt.java
+++ b/clients/src/main/java/org/apache/kafka/common/security/oauthbearer/internals/secured/SerializedJwt.java
@@ -44,7 +44,7 @@ public SerializedJwt(String token) {
String[] splits = token.split("\\.");
if (splits.length != 3)
- throw new ValidateException(String.format("Malformed JWT provided (%s); expected three sections (header, payload, and signature), but %s sections provided",
+ throw new ValidateException(String.format("Malformed JWT provided (%s); expected three sections (header, payload, and signature), but %d sections provided",
token, splits.length));
this.token = token.trim();
diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java
index 8360bf2b9ff99..1610c208f1a9c 100644
--- a/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java
+++ b/clients/src/test/java/org/apache/kafka/clients/producer/KafkaProducerTest.java
@@ -118,6 +118,7 @@
import static java.util.Collections.emptyMap;
import static java.util.Collections.singletonMap;
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
@@ -493,6 +494,12 @@ public void testNoSerializerProvided() {
Properties producerProps = new Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9000");
assertThrows(ConfigException.class, () -> new KafkaProducer(producerProps));
+
+ final Map configs = new HashMap<>();
+ configs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9999");
+
+ // Invalid value null for configuration key.serializer: must be non-null.
+ assertThrows(ConfigException.class, () -> new KafkaProducer(configs));
}
@Test
@@ -2399,4 +2406,23 @@ public KafkaProducer newKafkaProducer() {
}
}
+ @Test
+ void testDeliveryTimeoutAndLingerMsConfig() {
+ final Map configs = new HashMap<>();
+ configs.put(ProducerConfig.CLIENT_ID_CONFIG, "testDeliveryTimeoutAndLingerMsConfig");
+ configs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9999");
+ configs.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 1000);
+ configs.put(ProducerConfig.LINGER_MS_CONFIG, 1000);
+ configs.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 1);
+
+ // delivery.timeout.ms should be equal to or larger than linger.ms + request.timeout.ms
+ assertThrows(KafkaException.class, () -> new KafkaProducer<>(configs, new StringSerializer(), new StringSerializer()));
+
+ configs.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 1000);
+ configs.put(ProducerConfig.LINGER_MS_CONFIG, 999);
+ configs.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 1);
+
+ assertDoesNotThrow(() -> new KafkaProducer<>(configs, new StringSerializer(), new StringSerializer()));
+ }
+
}
diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/BuiltInPartitionerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/BuiltInPartitionerTest.java
index 734aedc483ad1..69546aaa6ed71 100644
--- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/BuiltInPartitionerTest.java
+++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/BuiltInPartitionerTest.java
@@ -29,8 +29,10 @@
import java.util.concurrent.atomic.AtomicInteger;
import static java.util.Arrays.asList;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class BuiltInPartitionerTest {
@@ -195,4 +197,10 @@ public void adaptivePartitionsTest() {
"Partition " + i + " was chosen " + frequencies[i] + " times");
}
}
+
+ @Test
+ void testStickyBatchSizeMoreThatZero() {
+ assertThrows(IllegalArgumentException.class, () -> new BuiltInPartitioner(logContext, TOPIC_A, 0));
+ assertDoesNotThrow(() -> new BuiltInPartitioner(logContext, TOPIC_A, 1));
+ }
}
diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java
index 4beeaaf876a74..6ea080a59a349 100644
--- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java
+++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/integration/MirrorConnectorsIntegrationBaseTest.java
@@ -262,8 +262,6 @@ public void shutdownClusters() throws Exception {
for (String x : backup.connectors()) {
backup.deleteConnector(x);
}
- deleteAllTopics(primary.kafka());
- deleteAllTopics(backup.kafka());
} finally {
shuttingDown = true;
try {
@@ -1049,17 +1047,6 @@ protected static void waitForTopicCreated(EmbeddedConnectCluster cluster, String
}
}
- /*
- * delete all topics of the input kafka cluster
- */
- private static void deleteAllTopics(EmbeddedKafkaCluster cluster) throws Exception {
- try (final Admin adminClient = cluster.createAdminClient()) {
- Set topicsToBeDeleted = adminClient.listTopics().names().get();
- log.debug("Deleting topics: {} ", topicsToBeDeleted);
- adminClient.deleteTopics(topicsToBeDeleted).all().get();
- }
- }
-
/*
* retrieve the config value based on the input cluster, topic and config name
*/
diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java
index 7f04914c8d12f..7bbcf0f686161 100644
--- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java
+++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java
@@ -311,7 +311,7 @@ protected Map> buildTasksConfig(String conn
Map> configs = new HashMap<>();
for (ConnectorTaskId cti : configState.tasks(connector)) {
- configs.put(cti, configState.taskConfig(cti));
+ configs.put(cti, configState.rawTaskConfig(cti));
}
return configs;
diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java
index c7fee9e671536..babb157f772c2 100644
--- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java
+++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java
@@ -579,7 +579,7 @@ public int hashCode() {
public void tasksConfig(String connName, Callback