This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 0d26e551292 KAFKA-20825 Fix native image startup with share group DLQ
manager (#22900)
0d26e551292 is described below
commit 0d26e55129204c22e9157b80e8f7d71f76cd05ca
Author: YANG-SYUAN CHOU <[email protected]>
AuthorDate: Tue Jul 28 18:31:29 2026 +0800
KAFKA-20825 Fix native image startup with share group DLQ manager (#22900)
Native Kafka images built from trunk can fail during broker startup with
a `ClassNotFoundException` for `DefaultShareGroupDLQManager`.
During startup, `BrokerServer.createShareGroupDLQManager` reads the
`group.share.dlq.manager.class.name` configuration. This configuration
defaults to
`org.apache.kafka.server.share.dlq.DefaultShareGroupDLQManager`, so
every broker process follows this code path even when the user has not
explicitly configured a DLQ manager.
The method previously passed the configured class name to
`Utils.loadClass` before comparing it against the two supported
implementations. Runtime class-name loading requires additional
reachability metadata in a GraalVM native image, causing the native
executable to fail before the broker could finish starting.
This reflective class loading is unnecessary because arbitrary DLQ
manager implementations are not supported. The method only accepts
`DefaultShareGroupDLQManager` and `NoOpShareGroupDLQManager`, and both
are instantiated directly.
This change compares the configured class name directly against the two
supported implementations. The behavior for the default, no-op, empty,
and unsupported configurations remains unchanged, while native-image
startup no longer depends on reflection metadata for the DLQ manager.
The original native-image failure was reported in:
https://github.com/apache/kafka/pull/22379#pullrequestreview-4731595296
Testing:
```text
./gradlew core:compileScala
```
Result:
```text
> Task :core:compileScala
BUILD SUCCESSFUL
```
A release archive and GraalVM native Docker image were then built from
the modified working tree:
```text
./gradlew clean releaseTarGz
docker build --no-cache --progress=plain ...
```
Native-image build result:
```text
[8/8] Creating image...
127.17MB in total
Produced artifacts: /app/kafka/kafka.Kafka (executable)
Finished generating 'kafka.Kafka' in 2m 32s.
```
The resulting native image was started as a single-node KRaft broker
using the default DLQ manager. The relevant output was collected with:
```bash
docker logs kafka-native-dlq-test 2>&1 |
grep -E \
"ClassNotFoundException|DefaultShareGroupDLQManager|Kafka Server
started"
```
Output:
```text
group.share.dlq.manager.class.name =
org.apache.kafka.server.share.dlq.DefaultShareGroupDLQManager
[KafkaRaftServer nodeId=1] Kafka Server started
```
No `ClassNotFoundException` or `NoClassDefFoundError` was present, and
the container remained running after startup.
Reviewers: Chia-Ping Tsai <[email protected]>, Gaurav Narula
<[email protected]>
---
core/src/main/scala/kafka/server/BrokerServer.scala | 8 ++++----
1 file changed, 4 insertions(+), 4 deletions(-)
diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala
b/core/src/main/scala/kafka/server/BrokerServer.scala
index f7f34c17537..b23ae34e1fe 100644
--- a/core/src/main/scala/kafka/server/BrokerServer.scala
+++ b/core/src/main/scala/kafka/server/BrokerServer.scala
@@ -774,9 +774,9 @@ class BrokerServer(
}
private def createShareGroupDLQManager(): ShareGroupDLQManager = {
- if (config.shareGroupConfig.shareGroupDLQManagerClassName.nonEmpty) {
- val klass =
Utils.loadClass(config.shareGroupConfig.shareGroupDLQManagerClassName,
classOf[Object]).asInstanceOf[Class[ShareGroupDLQManager]]
- if (klass.getName.equals(classOf[DefaultShareGroupDLQManager].getName)) {
+ val className = config.shareGroupConfig.shareGroupDLQManagerClassName
+ if (className.nonEmpty) {
+ if (className.equals(classOf[DefaultShareGroupDLQManager].getName)) {
DefaultShareGroupDLQManager.instance(
NetworkUtils.buildNetworkClient("ShareGroupDLQManager", config,
metrics, Time.SYSTEM, new LogContext(s"[ShareGroupDLQManager
broker=${config.brokerId}]")),
new ShareCoordinatorMetadataCacheHelperImpl(metadataCache, key =>
shareCoordinator.partitionFor(key), config.interBrokerListenerName,
groupConfigManager, () => config.messageMaxBytes),
@@ -785,7 +785,7 @@ class BrokerServer(
shareGroupMetrics,
shareGroupLogReader
)
- } else if
(klass.getName.equals(classOf[NoOpShareGroupDLQManager].getName)) {
+ } else if (className.equals(classOf[NoOpShareGroupDLQManager].getName)) {
info("Using no-op share group DLQ manager")
new NoOpShareGroupDLQManager()
} else {