Skip to content
Merged
Show file tree
Hide file tree
Changes from 37 commits
Commits
Show all changes
101 commits
Select commit Hold shift + click to select a range
e25f5ab
KAFKA-14589 WIP
nizhikov Sep 7, 2023
6da9b64
Merge branch 'trunk' into KAFKA-14589
nizhikov Oct 2, 2023
1bba7b5
Merge branch 'trunk' into KAFKA-14589
nizhikov Oct 22, 2023
b2022aa
KAFKA-14589 WIP
nizhikov Oct 22, 2023
0554aaa
KAFKA-14589 WIP
nizhikov Oct 22, 2023
304f584
KAFKA-14589 WIP
nizhikov Oct 22, 2023
a1e3d6c
Merge branch 'trunk' into KAFKA-14589
nizhikov Oct 24, 2023
d57e7c1
KAFKA-14589 WIP
nizhikov Oct 24, 2023
b650bb3
KAFKA-14589 WIP
nizhikov Oct 25, 2023
1f4ea90
KAFKA-14589 WIP
nizhikov Oct 26, 2023
7608683
Merge branch 'trunk' into KAFKA-14589
nizhikov Oct 26, 2023
00b5b01
KAFKA-14589 WIP
nizhikov Oct 26, 2023
6d52f02
KAFKA-14589 WIP
nizhikov Oct 26, 2023
b9c10b9
KAFKA-14589 Command implemented
nizhikov Oct 27, 2023
13e517f
KAFKA-14589 Codestyle fixes.
nizhikov Oct 27, 2023
cf9fcbd
KAFKA-14589 Codestyle fixes.
nizhikov Oct 27, 2023
daba249
Merge branch 'trunk' into KAFKA-14589
nizhikov Oct 30, 2023
ef64af8
KAFKA-14589 ConsumerGroupCommandTest added.
nizhikov Oct 31, 2023
3b740c9
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 1, 2023
4fc836c
KAFKA-14589 Merging with trunk
nizhikov Nov 1, 2023
a1b6398
KAFKA-14589 ListConsumerGroupTest rewritten in java
nizhikov Nov 5, 2023
64a3477
KAFKA-14589 DeleteOffsetsConsumerGroupCommandIntegrationTest rewritte…
nizhikov Nov 5, 2023
c31cba6
KAFKA-14589 DeleteConsumerGroupsTest rewritten in java
nizhikov Nov 5, 2023
f890648
KAFKA-14589 ResetConsumerGroupOffsetTest rewritten in java
nizhikov Nov 6, 2023
a80a8dd
KAFKA-14589 DescribeConsumerGroupTest rewritten in java
nizhikov Nov 8, 2023
f17135b
KAFKA-14589 DescribeConsumerGroupTest rewritten in java
nizhikov Nov 8, 2023
722aa86
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 8, 2023
c829fd0
KAFKA-14589 ConsumerGroupServiceTest rewritten in java
nizhikov Nov 8, 2023
67ddbf1
KAFKA-14589 Transfer final tests and remove scala version of command
nizhikov Nov 8, 2023
303e592
KAFKA-14589 Fix scala 2.12 build
nizhikov Nov 8, 2023
dc07aa3
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 15, 2023
adcb688
KAFKA-14589 Checkstyle fix
nizhikov Nov 15, 2023
e671974
KAFKA-14589 Checkstyle fix
nizhikov Nov 15, 2023
457af5c
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 20, 2023
0efe53b
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 24, 2023
ad6564f
Merge branch 'trunk' into KAFKA-14589
nizhikov Nov 28, 2023
9191f6c
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 5, 2023
a837e84
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 6, 2023
fb5c789
KAFKA-14588 Code review fix
nizhikov Dec 6, 2023
38e43d0
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 12, 2023
c74f4d4
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 21, 2023
8c324ce
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 23, 2023
377270a
Merge branch 'trunk' into KAFKA-14589
nizhikov Dec 27, 2023
8063da1
Merge branch 'trunk' into KAFKA-14589
Jan 22, 2024
0e3f4f4
KAFKA-14589 Update to latest trunk changes
Jan 22, 2024
fc288fe
KAFKA-14589 Change package
Jan 22, 2024
73e2fba
Merge branch 'trunk' into KAFKA-14589
Jan 23, 2024
74cc1d7
KAFKA-14589 WIP
Jan 23, 2024
a7080ae
KAFKA-14589 WIP
Jan 23, 2024
26a437d
Merge branch 'trunk' into KAFKA-14589
Jan 24, 2024
5e49029
Merge branch 'trunk' into KAFKA-14589
Jan 26, 2024
700a358
Merge branch 'trunk' into KAFKA-14589
Jan 26, 2024
1f4b547
KAFKA-14589 Merge fixes
Jan 26, 2024
85cd1d3
KAFKA-14589 Correct package
Jan 29, 2024
96ff15f
Merge branch 'trunk' into KAFKA-14589
Jan 29, 2024
9ae4e2a
KAFKA-14589 Reflection changes from #15211
Jan 29, 2024
a6b1a0b
Merge branch 'trunk' into KAFKA-14589
Feb 1, 2024
6307859
KAFKA-14589 Reflecting changes of 6c09cc9586f823dff32c96131f6f377afae…
Feb 1, 2024
6c870d4
Merge branch 'trunk' into KAFKA-14589
Feb 5, 2024
d29d611
Merge branch 'trunk' into KAFKA-14589
Feb 6, 2024
04f80f5
Merge branch 'trunk' into KAFKA-14589
Feb 13, 2024
170b800
KAFKA-14589 Fixes after merge
Feb 13, 2024
2030962
KAFKA-14589 Fixes after merge
Feb 13, 2024
3569971
KAFKA-14589 Fixes after merge
Feb 13, 2024
e4e0a1e
Merge branch 'trunk' into KAFKA-14589
Feb 13, 2024
b74bd03
KAFKA-14589 Reduce tests run
Feb 13, 2024
a34c6c6
Merge branch 'trunk' into KAFKA-14589
Feb 13, 2024
8426c04
Merge branch 'trunk' into KAFKA-14589
Feb 13, 2024
fd29b3f
KAFKA-14589 Tests rewritten
Feb 14, 2024
54f1017
KAFKA-14589 Reduce changes
Feb 14, 2024
4c60e22
KAFKA-14589 Reduce changes
Feb 14, 2024
a1254d8
Merge branch 'trunk' into KAFKA-14589
Feb 15, 2024
341c81b
Merge branch 'trunk' into KAFKA-14589
Feb 20, 2024
419c0ea
KAFKA-14588 WIP
Feb 20, 2024
4660f76
Merge branch 'trunk' into KAFKA-14589
Feb 27, 2024
3b1d38b
Merge branch 'trunk' into KAFKA-14589
Feb 29, 2024
b4539d8
KAFKA-14588 Reflect changes from KAFKA-15462
Feb 29, 2024
8a173f2
KAFKA-14588 Reflect changes from KAFKA-15462
Feb 29, 2024
311e9fa
Merge branch 'trunk' into KAFKA-14589
Mar 4, 2024
e406967
Merge branch 'trunk' into KAFKA-14589
Mar 5, 2024
3933267
Merge branch 'trunk' into KAFKA-14589
Mar 6, 2024
837a9b2
KAFKA-14588 Code review fixes
Mar 6, 2024
3ade3dc
KAFKA-14589 Revert unnecessary changes
Mar 6, 2024
e9018cd
Merge branch 'trunk' into KAFKA-14589
Mar 7, 2024
dfadfb7
KAFKA-14589 Cleanup
Mar 7, 2024
643dd26
KAFKA-14589 Cleanup
Mar 7, 2024
bfbdfa5
KAFKA-14589 Cleanup
Mar 7, 2024
95ec460
Merge branch 'trunk' into KAFKA-14589
Mar 8, 2024
8365091
KAFKA-14589 Merge updates
Mar 8, 2024
8560d5d
KAFKA-14589 Merge updates
Mar 8, 2024
528c56f
KAFKA-14589 Merge updates
Mar 8, 2024
0dd7e43
KAFKA-14589 Merge updates
Mar 8, 2024
7dbc9e4
Merge branch 'trunk' into KAFKA-14589
Mar 11, 2024
e1175a6
Merge branch 'trunk' into KAFKA-14589
Mar 12, 2024
c7600b8
KAFKA-14589 Code review changes
Mar 12, 2024
4f7befa
KAFKA-14589 Code review changes
Mar 12, 2024
3430af9
Merge branch 'trunk' into KAFKA-14589
Mar 15, 2024
23a01af
KAFKA-14589 Code review changes
Mar 15, 2024
ba9bf6a
Merge branch 'trunk' into KAFKA-14589
Mar 19, 2024
519f016
KAFKA-14589 Code review changes
Mar 19, 2024
6645d62
KAFKA-14589 Code review changes
Mar 19, 2024
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
2 changes: 1 addition & 1 deletion bin/kafka-consumer-groups.sh
Original file line number Diff line number Diff line change
Expand Up @@ -14,4 +14,4 @@
# See the License for the specific language governing permissions and
# limitations under the License.

exec $(dirname $0)/kafka-run-class.sh kafka.admin.ConsumerGroupCommand "$@"
exec $(dirname $0)/kafka-run-class.sh org.apache.kafka.tools.consumergroup.ConsumerGroupCommand "$@"
2 changes: 1 addition & 1 deletion bin/windows/kafka-consumer-groups.bat
Original file line number Diff line number Diff line change
Expand Up @@ -14,4 +14,4 @@ rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
rem See the License for the specific language governing permissions and
rem limitations under the License.

"%~dp0kafka-run-class.bat" kafka.admin.ConsumerGroupCommand %*
"%~dp0kafka-run-class.bat" org.apache.kafka.tools.consumergroup.ConsumerGroupCommand %*
5 changes: 5 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -1972,6 +1972,9 @@ project(':tools') {
implementation project(':log4j-appender')
implementation project(':tools:tools-api')
implementation libs.argparse4j
implementation libs.jacksonDatabind
implementation libs.jacksonDataformatCsv
implementation libs.jacksonJDK8Datatypes
implementation libs.slf4jApi
implementation libs.log4j
implementation libs.joptSimple
Expand All @@ -1986,6 +1989,8 @@ project(':tools') {
testImplementation project(':server-common')
testImplementation project(':server-common').sourceSets.test.output
testImplementation project(':connect:api')
testImplementation project(':metadata')
testImplementation project(':metadata').sourceSets.test.output
testImplementation project(':connect:runtime')
testImplementation project(':connect:runtime').sourceSets.test.output
testImplementation project(':storage:storage-api').sourceSets.main.output
Expand Down
6 changes: 6 additions & 0 deletions checkstyle/import-control.xml
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,12 @@
<allow pkg="scala" />
</subpackage>

<subpackage name="consumergroup">
<allow pkg="kafka.security"/>
<allow pkg="org.apache.kafka.metadata.authorizer"/>
<allow pkg="org.apache.kafka.tools"/>
</subpackage>

<subpackage name="other">
<allow pkg="org.apache.kafka.tools.reassign"/>
<allow pkg="kafka.log" />
Expand Down
10 changes: 10 additions & 0 deletions clients/src/main/java/org/apache/kafka/common/utils/Utils.java
Original file line number Diff line number Diff line change
Expand Up @@ -593,6 +593,16 @@ public static <T> String join(T[] strs, String separator) {
return join(Arrays.asList(strs), separator);
}

/**
* Create a string representation of a collection joined by ", ".
* @param collection The list of items
* @return The string representation.
*/
public static <T> String join(Collection<T> collection) {
Objects.requireNonNull(collection);
return mkString(collection.stream(), "", "", ", ");
}

/**
* Create a string representation of a collection joined by the given separator
* @param collection The list of items
Expand Down
1,165 changes: 0 additions & 1,165 deletions core/src/main/scala/kafka/admin/ConsumerGroupCommand.scala

This file was deleted.

2 changes: 1 addition & 1 deletion core/src/main/scala/kafka/utils/ToolsUtils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ object ToolsUtils {
/**
* This is a simple wrapper around `CommandLineUtils.printUsageAndExit`.
* It is needed for tools migration (KAFKA-14525), as there is no Java equivalent for return type `Nothing`.
* Can be removed once [[kafka.admin.ConsumerGroupCommand]], [[kafka.tools.ConsoleConsumer]]
* Can be removed once [[kafka.tools.ConsoleConsumer]]
* and [[kafka.tools.ConsoleProducer]] are migrated.
*
* @param parser Command line options parser.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ import java.util
import java.util.concurrent.ExecutionException
import java.util.regex.Pattern
import java.util.{Collections, Optional, Properties}
import kafka.admin.ConsumerGroupCommand.{ConsumerGroupCommandOptions, ConsumerGroupService}
import kafka.security.authorizer.{AclAuthorizer, AclEntry}
import kafka.security.authorizer.AclEntry.WildcardHost
import kafka.server.{BaseRequestTest, KafkaConfig}
Expand Down Expand Up @@ -1709,20 +1708,6 @@ class AuthorizerIntegrationTest extends BaseRequestTest {
createAdminClient().describeConsumerGroups(Seq(group).asJava).describedGroups().get(group).get()
}

@ParameterizedTest(name = TestInfoUtils.TestWithParameterizedQuorumName)
@ValueSource(strings = Array("zk", "kraft"))
def testDescribeGroupCliWithGroupDescribe(quorum: String): Unit = {
createTopicWithBrokerPrincipal(topic)
addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, DESCRIBE, ALLOW)), groupResource)
addAndVerifyAcls(Set(new AccessControlEntry(clientPrincipalString, WildcardHost, DESCRIBE, ALLOW)), topicResource)

val cgcArgs = Array("--bootstrap-server", bootstrapServers(), "--describe", "--group", group)
val opts = new ConsumerGroupCommandOptions(cgcArgs)
val consumerGroupService = new ConsumerGroupService(opts)
consumerGroupService.describeGroups()
consumerGroupService.close()
}

@ParameterizedTest(name = TestInfoUtils.TestWithParameterizedQuorumName)
@ValueSource(strings = Array("zk", "kraft"))
def testListGroupApiWithAndWithoutListGroupAcls(quorum: String): Unit = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ import org.apache.kafka.common.{KafkaException, TopicPartition}
import org.apache.kafka.common.errors.SaslAuthenticationException
import org.junit.jupiter.api.{AfterEach, BeforeEach, Test, TestInfo}
import org.junit.jupiter.api.Assertions._
import kafka.admin.ConsumerGroupCommand.{ConsumerGroupCommandOptions, ConsumerGroupService}
import kafka.server.KafkaConfig
import kafka.utils.{JaasTestUtils, TestUtils}
import kafka.zk.ConfigEntityChangeNotificationZNode
Expand All @@ -32,7 +31,7 @@ import org.junit.jupiter.params.ParameterizedTest
import org.junit.jupiter.params.provider.ValueSource

class SaslClientsWithInvalidCredentialsTest extends IntegrationTestHarness with SaslSetup {
private val kafkaClientSaslMechanism = "SCRAM-SHA-256"
val kafkaClientSaslMechanism = "SCRAM-SHA-256"
private val kafkaServerSaslMechanisms = List(kafkaClientSaslMechanism)
override protected val securityProtocol = SecurityProtocol.SASL_PLAINTEXT
override protected val serverSaslProperties = Some(kafkaServerSaslProperties(kafkaServerSaslMechanisms, kafkaClientSaslMechanism))
Expand Down Expand Up @@ -166,46 +165,7 @@ class SaslClientsWithInvalidCredentialsTest extends IntegrationTestHarness with
}
}

@Test
def testConsumerGroupServiceWithAuthenticationFailure(): Unit = {
val consumerGroupService: ConsumerGroupService = prepareConsumerGroupService

val consumer = createConsumer()
try {
consumer.subscribe(List(topic).asJava)

verifyAuthenticationException(consumerGroupService.listGroups())
} finally consumerGroupService.close()
}

@Test
def testConsumerGroupServiceWithAuthenticationSuccess(): Unit = {
createClientCredential()
val consumerGroupService: ConsumerGroupService = prepareConsumerGroupService

val consumer = createConsumer()
try {
consumer.subscribe(List(topic).asJava)

verifyWithRetry(consumer.poll(Duration.ofMillis(1000)))
assertEquals(1, consumerGroupService.listConsumerGroups().size)
}
finally consumerGroupService.close()
}

private def prepareConsumerGroupService = {
val propsFile = TestUtils.tempPropertiesFile(Map("security.protocol" -> "SASL_PLAINTEXT", "sasl.mechanism" -> kafkaClientSaslMechanism))

val cgcArgs = Array("--bootstrap-server", bootstrapServers(),
"--describe",
"--group", "test.group",
"--command-config", propsFile.getAbsolutePath)
val opts = new ConsumerGroupCommandOptions(cgcArgs)
val consumerGroupService = new ConsumerGroupService(opts)
consumerGroupService
}

private def createClientCredential(): Unit = {
def createClientCredential(): Unit = {
createScramCredentialsViaPrivilegedAdminClient(JaasTestUtils.KafkaScramUser2, JaasTestUtils.KafkaScramPassword2)
}

Expand All @@ -221,14 +181,14 @@ class SaslClientsWithInvalidCredentialsTest extends IntegrationTestHarness with
}
}

private def verifyAuthenticationException(action: => Unit): Unit = {
def verifyAuthenticationException(action: => Unit): Unit = {
val startMs = System.currentTimeMillis
assertThrows(classOf[Exception], () => action)
val elapsedMs = System.currentTimeMillis - startMs
assertTrue(elapsedMs <= 5000, s"Poll took too long, elapsed=$elapsedMs")
}

private def verifyWithRetry(action: => Unit): Unit = {
def verifyWithRetry(action: => Unit): Unit = {
var attempts = 0
TestUtils.waitUntilTrue(() => {
try {
Expand Down
211 changes: 0 additions & 211 deletions core/src/test/scala/unit/kafka/admin/ConsumerGroupCommandTest.scala

This file was deleted.

Loading