confluentinc / confluentinc/parallel-consumer

Migration from java 8 to 21, from springboot 2.7 to 3.5

Open
#878 3 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
299
Forks
172
PR merge metrics
No merged PRs in 30d

Description

Hello,
My teams is migrating from java 8 to 21 and also springboot 3.7.x to 3.5 and I have errors.

I also upgrate the libraries
from
7.2.1
to
7.9.0

and
from
0.5.2.1
to
0.5.3.2

I don't understand why I have this problem.

Thank fo yout help

My pom.xml
`

4.0.0

org.springframework.boot
spring-boot-starter-parent
3.5.0


xxx
notif_service
jar
v2025.1.1-snapshoot


3.5.0
21
21
21


7.9.0
0.5.2.8
4.9.10
UTF-8
UTF-8
0.10.6

0.8.13
jacoco
reuseReports
${project.build.directory}/site/jacoco/jacoco.xml
java
src/main/java
src/test
**/src/test/**
-Duser.timezone=UTC



org.springframework.boot
spring-boot-starter


org.springframework.boot
spring-boot-starter-web


org.springframework.boot
spring-boot-starter-data-jpa


org.apache.velocity
velocity-engine-core
2.3


de.codecentric
spring-boot-admin-starter-client
3.4.5


io.confluent
kafka-streams-avro-serde
${confluent.version}


io.confluent.parallelconsumer
parallel-consumer-core
${parallel-consumer.version}


org.springframework
spring-context


org.springframework
spring-core


org.springframework
spring-jdbc


com.microsoft.sqlserver
mssql-jdbc
runtime


io.vavr
vavr
${vavr.version}



org.junit.jupiter
junit-jupiter
test


org.springframework.boot
spring-boot-starter-test
test


org.hsqldb
hsqldb
test



org.springframework.boot
spring-boot-starter-actuator



io.micrometer
micrometer-registry-prometheus


org.springframework.boot
spring-boot-starter-logging


org.slf4j
jul-to-slf4j


org.apache.logging.log4j
log4j-to-slf4j




org.projectlombok
lombok


kur-notification-service


pl.project13.maven
git-commit-id-plugin
${git-commit-id-plugin.version}


org.springframework.boot
spring-boot-maven-plugin
${spring-boot.version}


build-info

build-info



${java.version}
${kafka.version}






org.jacoco
jacoco-maven-plugin
${jacoco.version}


jacoco-initialize

prepare-agent




jacoco-site
test

report





org.apache.maven.plugins
maven-assembly-plugin


false

src/main/assembly/assembly_notification_service.xml




package
package

single





org.apache.maven.plugins
maven-surefire-plugin



`

My code :
`import io.confluent.kafka.serializers.AbstractKafkaSchemaSerDeConfig;
import io.confluent.kafka.serializers.subject.TopicNameStrategy;
import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde;
import org.apache.avro.specific.SpecificRecord;

import java.util.HashMap;
import java.util.Map;
import java.util.Properties;
public class SerdesUtil {

public static SpecificAvroSerde getValueSerdes(Class type, Properties props, boolean isKey) {
Map registry = getPropertiesMap(props);
registry.put(AbstractKafkaSchemaSerDeConfig.VALUE_SUBJECT_NAME_STRATEGY, TopicNameStrategy.class.getName());
SpecificAvroSerde dataAvroSerde = new SpecificAvroSerde<>();
dataAvroSerde.configure(registry, isKey);
return dataAvroSerde;
}

private static Map getPropertiesMap(Properties props) {
Map originals = new HashMap<>();
for (final String name : props.stringPropertyNames()) {
originals.put(name, String.valueOf(props.get(name)));
}
return originals;
}
}`

`import io.confluent.parallelconsumer.ParallelConsumerOptions;
import io.confluent.parallelconsumer.ParallelStreamProcessor;
import org.apache.avro.specific.SpecificRecord;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.common.serialization.Serdes;

import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Collections;
import java.util.Properties;

import static io.confluent.parallelconsumer.ParallelConsumerOptions.ProcessingOrder.KEY;
import static java.lang.Integer.parseInt;
import static org.apache.kafka.clients.producer.ProducerConfig.TRANSACTIONAL_ID_CONFIG;

public class GenericParallelConsumer {

private final Properties properties;
private final String inputTopic;

private static final String PROPERTY_INPUT_TOPIC_PREFIX = "kafka.topic.input.";
private static final String PROPERTY_MAX_CONCURRENCY = "maxConcurrency";

public GenericParallelConsumer(Properties properties, String inputTopic) {
this.properties = properties;
this.inputTopic = inputTopic;
}

public ParallelStreamProcessor createParallelConsumer(String groupId) {
try {
properties.put(TRANSACTIONAL_ID_CONFIG, properties.getProperty(groupId) + "-" + InetAddress.getLocalHost().getHostName());
} catch (UnknownHostException e) {
throw new RuntimeException(e);
}

Consumer consumer = new KafkaConsumer<>(properties, null, SerdesUtil.getValueSerdes(SpecificRecord.class, properties, false).deserializer());
Producer producer = new KafkaProducer<>(properties, Serdes.String().serializer(), SerdesUtil.getValueSerdes(SpecificRecord.class, properties, false).serializer());
ParallelConsumerOptions options =
ParallelConsumerOptions.builder()
.ordering(KEY)
.maxConcurrency(parseInt(properties.getProperty(PROPERTY_MAX_CONCURRENCY)))
.consumer(consumer)
.producer(producer)
.commitMode(ParallelConsumerOptions.CommitMode.PERIODIC_TRANSACTIONAL_PRODUCER)
.build();

ParallelStreamProcessor parallelConsumer =
ParallelStreamProcessor.createEosStreamProcessor(options);

parallelConsumer.subscribe(Collections.singleton(properties.getProperty(PROPERTY_INPUT_TOPIC_PREFIX.concat(inputTopic))));
return parallelConsumer;
}
}`

The values of my options :
`
"compression.type" -> "zstd"

"value.deserializer" -> "io.confluent.kafka.serializers.KafkaAvroDeserializer"

"group.id.email" -> "EmailPreparationService"

"auto.register.schemas" -> "false"

"group.id.report" -> "EmailPreparationServiceReport"

"group.id" -> "EmailPreparationServiceReport"

"kafka.topic.input.report" -> "KureReportRequested_V1"

"bootstrap.servers" -> "server:8080"

"schema.registry.ssl.truststore.location" -> "C:/Dev/Configs/CertifKure/kafka.client.truststore.jks"

"kafka.topic.input.notification" -> "KureNotificationPrepared_V1"

"maxConcurrency" -> "1"

"transactional.id" -> "EmailPreparationServiceReport-L0009228"

"schema.registry.url" -> "https://schema.com"

"outputTopic" -> "KureEmailPrepared_V1"

"enable.auto.commit" -> "false"

"sasl.mechanism" -> "SCRAM-SHA-512"

"schema.registry.basic.auth.credentials.source" -> "USER_INFO"

"sasl.jaas.config" -> "org.apache.kafka.common.security.scram.ScramLoginModule required username="user" password="pwd";"

"group.id.reminder" -> "eminderPreparationService"

"ssl.truststore.password" -> "pwd"

"ssl.endpoint.identification.algorithm" -> "https"

"key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer"

"kafka.topic.input.reminder" -> "CommentPublished_V1"

"schema.registry.ssl.truststore.password" -> "pwd"

"errorTopic" -> "GlobalError_DLQ_V1"

"security.protocol" -> "SASL_SSL"

"avro.remove.java.properties" -> "true"

"ssl.truststore.location" -> "C:/Dev/Configs/CertifKure/kafka.client.truststore.jks"

"isolation.level" -> "read_committed"

"schema.registry.basic.auth.user.info" -> "user:pwd
`

The error messages :
`16-06-2025 10:48:50.147 [Thread-3] INFO org.apache.kafka.common.utils.AppInfoParser. - Kafka version: 3.9.1
16-06-2025 10:48:50.147 [Thread-3] INFO org.apache.kafka.common.utils.AppInfoParser. - Kafka commitId: f745dfdcee2b9851
16-06-2025 10:48:50.147 [Thread-3] INFO org.apache.kafka.common.utils.AppInfoParser. - Kafka startTimeMs: 1750063730146
16-06-2025 10:48:50.156 [Thread-3] WARN io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.getAutoCommitEnabled - Encountered unknown consumer delegate class org.apache.kafka.clients.consumer.KafkaConsumer
16-06-2025 10:48:50.156 [Thread-1] WARN io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.getAutoCommitEnabled - Encountered unknown consumer delegate class org.apache.kafka.clients.consumer.KafkaConsumer
16-06-2025 10:48:50.156 [Thread-2] WARN io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.getAutoCommitEnabled - Encountered unknown consumer delegate class org.apache.kafka.clients.consumer.KafkaConsumer
Exception in thread "Thread-3" Exception in thread "Thread-1" Exception in thread "Thread-2" io.confluent.parallelconsumer.ParallelConsumerException: Unable to check whether auto commit is enabled for consumer type class org.apache.kafka.clients.consumer.KafkaConsumer. This exception can be ignored by enabling the ignoreReflectiveAccessExceptionsForAutoCommitDisabledCheck option.
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.checkAutoCommitIsDisabled(AbstractParallelEoSStreamProcessor.java:501)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.validateConfiguration(AbstractParallelEoSStreamProcessor.java:336)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:289)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:272)
at io.confluent.parallelconsumer.ParallelEoSStreamProcessor.(ParallelEoSStreamProcessor.java:46)
at io.confluent.parallelconsumer.ParallelStreamProcessor.createEosStreamProcessor(ParallelStreamProcessor.java:26)
at xxx.service.common.consumer.GenericParallelConsumer.createParallelConsumer(GenericParallelConsumer.java:54)
at xxx.service.reportpreparation.ReportPreparationServiceExecutor.run(ReportPreparationServiceExecutor.java:40)
at java.base/java.lang.Thread.run(Thread.java:1583)
io.confluent.parallelconsumer.ParallelConsumerException: Unable to check whether auto commit is enabled for consumer type class org.apache.kafka.clients.consumer.KafkaConsumer. This exception can be ignored by enabling the ignoreReflectiveAccessExceptionsForAutoCommitDisabledCheck option.
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.checkAutoCommitIsDisabled(AbstractParallelEoSStreamProcessor.java:501)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.validateConfiguration(AbstractParallelEoSStreamProcessor.java:336)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:289)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:272)
at io.confluent.parallelconsumer.ParallelEoSStreamProcessor.(ParallelEoSStreamProcessor.java:46)
at io.confluent.parallelconsumer.ParallelStreamProcessor.createEosStreamProcessor(ParallelStreamProcessor.java:26)
at xxx.service.common.consumer.GenericParallelConsumer.createParallelConsumer(GenericParallelConsumer.java:54)
at xxx.service.reminder.ReminderPreparationServiceExecutor.run(ReminderPreparationServiceExecutor.java:44)
at java.base/java.lang.Thread.run(Thread.java:1583)
io.confluent.parallelconsumer.ParallelConsumerException: Unable to check whether auto commit is enabled for consumer type class org.apache.kafka.clients.consumer.KafkaConsumer. This exception can be ignored by enabling the ignoreReflectiveAccessExceptionsForAutoCommitDisabledCheck option.
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.checkAutoCommitIsDisabled(AbstractParallelEoSStreamProcessor.java:501)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.validateConfiguration(AbstractParallelEoSStreamProcessor.java:336)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:289)
at io.confluent.parallelconsumer.internal.AbstractParallelEoSStreamProcessor.(AbstractParallelEoSStreamProcessor.java:272)
at io.confluent.parallelconsumer.ParallelEoSStreamProcessor.(ParallelEoSStreamProcessor.java:46)
at io.confluent.parallelconsumer.ParallelStreamProcessor.createEosStreamProcessor(ParallelStreamProcessor.java:26)
at xxx.service.common.consumer.GenericParallelConsumer.createParallelConsumer(GenericParallelConsumer.java:54)
at xxx.service.email.EmailPreparationServiceExecutor.run(EmailPreparationServiceExecutor.java:56)
at java.base/java.lang.Thread.run(Thread.java:1583)`

Contributor guide

No contributing guide indexed for this repository

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.