Oops, the picture of the application.yml file is t...
# all-things-deployment
p
Oops, the picture of the application.yml file is transcribed, so it looks like there is a typo. Actually there are no typos.
Copy code
kafka:
  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.
Copy code
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;
  }
}