From 11c3acdf2ce1632fab53cd94df51ac08d57e8199 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 15 Mar 2022 13:48:07 -0400 Subject: [PATCH] Remove usage of isOnlyLogRecordMetadata The `ConsumerProperties.isOnlyLogRecordMetadata` has been removed in Spring for Apache Kafka in favor of `KafkaUtils.setConsumerRecordFormatter()` --- .../integration/kafka/inbound/KafkaMessageSource.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java index c4623e8ff7..0f4720e8d0 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageSource.java @@ -691,8 +691,6 @@ public class KafkaMessageSource extends AbstractMessageSource impl private final boolean isSyncCommits; - private final boolean logOnlyMetadata; - private volatile boolean acknowledged; private boolean autoAckEnabled = true; @@ -716,7 +714,6 @@ public class KafkaMessageSource extends AbstractMessageSource impl consumerProperties != null ? consumerProperties.getCommitLogLevel() : LogIfLevelEnabled.Level.DEBUG); - this.logOnlyMetadata = consumerProperties != null && consumerProperties.isOnlyLogRecordMetadata(); } @Override @@ -761,7 +758,6 @@ public class KafkaMessageSource extends AbstractMessageSource impl }) .collect(Collectors.toList()); if (rewound.size() > 0) { - KafkaUtils.setLogOnlyMetadata(this.logOnlyMetadata); this.logger.warn(() -> "Rolled back " + KafkaUtils.format(record) + " later in-flight offsets " + rewound + " will also be re-fetched"); @@ -771,7 +767,6 @@ public class KafkaMessageSource extends AbstractMessageSource impl } private void commitIfPossible(ConsumerRecord record) { // NOSONAR - KafkaUtils.setLogOnlyMetadata(this.logOnlyMetadata); if (this.ackInfo.isRolledBack()) { this.logger.warn(() -> "Cannot commit offset for " + KafkaUtils.format(record)