Kafka Exception: The group member needs to have a valid member id before actually entering a consumer group

Viewed 23

I have this consumer that appears to be connected to the kafka topic but after seeing the logs I'm seeing that the consumer Join group failed with org.apache.kafka.common.errors.MemberIdRequiredException: The group member needs to have a valid member id before actually entering a consumer group. I've tried looking at a few other setups online and I don't see anyone explicitly setting a member ID. I don't see this exception thrown locally so I'm hoping that maybe I just need to do some fine-tuning on the consumer configurations but I'm not sure what needs to be changed.

Kafka Consumer

@KafkaListener(
        topics = ["my.topic"],
        groupId ="kafka-consumer-group"
    )
    fun consume(consumerRecord: ConsumerRecord<String, Record>, ack: Acknowledgment) {
        val value = consumerRecord.value()
        val eventKey = consumerRecord.key()
        val productNotificationEvent = buildProductNotificationEvent(value)
        val consumerIdsByEventKey = eventKey.let { favoritesRepository.getConsumerIdsByEntityId(it) }        
        ack.acknowledge()
}

Kafka Consumer Config

@EnableKafka
@Configuration
class KafkaConsumerConfig(
    @Value("\${KAFKA_SERVER}")
    private val kafkaServer: String,
    @Value("\${KAFKA_SASL_USERNAME}")
    private val userName: String,
    @Value("\${KAFKA_SASL_PASSWORD}")
    private val password: String,
    @Value("\${KAFKA_SCHEMA_REGISTRY_URL}")
    private val schemaRegistryUrl: String,
    @Value("\${KAFKA_SCHEMA_REGISTRY_USER_INFO}")
    private val schemaRegistryUserInfo: String,
    @Value("\${KAFKA_CONSUMER_MAX_POLL_IN_MS}")
    private val maxPollIntervalMsConfig: Int,
    @Value("\${KAFKA_CONSUMER_MAX_POLL_RECORDS}")
    private val maxPollRecords: Int
) {

    @Bean
    fun consumerFactory(): ConsumerFactory<String?, Any?> {
        val props: MutableMap<String, Any> = mutableMapOf()
        props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = kafkaServer
        props[ConsumerConfig.GROUP_ID_CONFIG] = KAFKA_CONSUMER_GROUP_ID
        props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = StringDeserializer::class.java
        props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = KafkaAvroDeserializer::class.java
        props[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = "earliest"
        props[ConsumerConfig.MAX_POLL_RECORDS_CONFIG] = maxPollRecords
        props[ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG] = maxPollIntervalMsConfig

        props[CommonClientConfigs.SECURITY_PROTOCOL_CONFIG] = "SASL_SSL"
        props[SaslConfigs.SASL_MECHANISM] = "PLAIN"
        val module = "org.apache.kafka.common.security.plain.PlainLoginModule"
        val jaasConfig = String.format(
            "%s required username=\"%s\" password=\"%s\";",
            module,
            userName,
            password
        )
        props[SaslConfigs.SASL_JAAS_CONFIG] = jaasConfig

        props[KafkaAvroDeserializerConfig.VALUE_SUBJECT_NAME_STRATEGY] = TopicRecordNameStrategy::class.java
        props[KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG] = schemaRegistryUrl
        props[KafkaAvroDeserializerConfig.BASIC_AUTH_CREDENTIALS_SOURCE] = "USER_INFO"
        props[KafkaAvroDeserializerConfig.USER_INFO_CONFIG] = schemaRegistryUserInfo

        return DefaultKafkaConsumerFactory(props)
    }

    @Bean
    fun kafkaListenerContainerFactory(): ConcurrentKafkaListenerContainerFactory<String, Any>? {
        val factory = ConcurrentKafkaListenerContainerFactory<String, Any>()
        factory.consumerFactory = consumerFactory()
        factory.containerProperties.ackMode = ContainerProperties.AckMode.MANUAL_IMMEDIATE
        factory.containerProperties.isSyncCommits = true
        return factory
    }

    companion object {
        private const val KAFKA_CONSUMER_GROUP_ID = "kafka-consumer-group"
    }
}
0 Answers
Related