polite-actor-701
01/26/2023, 12:54 PMkafka:
listener:
concurrency: ${KAFKA_LISTENER_CONCURRENCY:1}
bootstrapServers: ${KAFKA_BOOTSTRAP_SERVER:<http://localhost:9092>}
schemaRegistry:
type: ${SCHEMA_REGISTRY_TYPE:KAFKA} # KAFKA or AWS_GLUE
url: ${KAFKA_SCHEMAREGISTRY_URL:<http://localhost:8081>} # Application only for type = kafka
awsGlue:
region: ${AWS_GLUE_SCHEMA_REGISTRY_REGION:us-east-1}
registryName: ${AWS_GLUE_SCHEMA_REGISTRY_NAME:#{null}}
consumer:
maxPollRecords: ${KAFKA_CONSUMER_MAX_POLL_RECORDS:200}
I also tried reducing KAFKA_CONSUMER_MAX_POLL_RECORDS to 1, but got the same error.
And I haven't modified anything related except for the two files in the picture above and adding env when running the container.
I'll insert the entire code for your reference.
package com.linkedin.gms.factory.kafka;
import java.time.Duration;
import java.util.Arrays;
import lombok.extern.slf4j.Slf4j;
import org.apache.avro.generic.GenericRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerContainerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
@Slf4j
@Configuration
@EnableConfigurationProperties(KafkaProperties.class)
public class SimpleKafkaConsumerFactory {
@Value("${KAFKA_BOOTSTRAP_SERVER:<http://localhost:9092>}")
private String kafkaBootstrapServers;
@Value("${KAFKA_CONSUMER_MAX_POLL_RECORDS:200}")
private Integer maxPollRecords;
@Bean(name = "simpleKafkaConsumer")
protected KafkaListenerContainerFactory<?> createInstance(KafkaProperties properties) {
KafkaProperties.Consumer consumerProps = properties.getConsumer();
// Specify (de)serializers for record keys and for record values.
consumerProps.setKeyDeserializer(StringDeserializer.class);
consumerProps.setValueDeserializer(StringDeserializer.class);
// Records will be flushed every 10 seconds.
consumerProps.setEnableAutoCommit(true);
consumerProps.setAutoCommitInterval(Duration.ofSeconds(10));
// [KB_SEOHEE]
consumerProps.setMaxPollRecords(maxPollRecords);
System.out.println("---------------------------------------"+consumerProps.getMaxPollRecords());
// KAFKA_BOOTSTRAP_SERVER has precedence over SPRING_KAFKA_BOOTSTRAP_SERVERS
if (kafkaBootstrapServers != null && kafkaBootstrapServers.length() > 0) {
consumerProps.setBootstrapServers(Arrays.asList(kafkaBootstrapServers.split(",")));
} // else we rely on KafkaProperties which defaults to localhost:9092
ConcurrentKafkaListenerContainerFactory<String, GenericRecord> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(properties.buildConsumerProperties()));
<http://log.info|log.info>("Simple KafkaListenerContainerFactory built successfully");
return factory;
}
}