Я пытаюсь создать Group.id Динамически, как показано ниже. Цель этого состоит в том, чтобы употреблять сообщение для всего запуска экземпляра моего приложения. Например, есть 15 экземпляров, работающих с различным именем хоста, я хочу создать 15 группового идентификатора, так что каждое сообщение в теме потребляется всеми группами потребителей. фрагмент < /p>
@KafkaListener(
id = BeanConstants.CACHE_UPDATE_REQUEST_LISTENER_ID,
groupId = "#{ 'test-cache-update-event-cg-' + T(java.lang.System).getenv('HOSTNAME')}")
public void consume(ConsumerRecord cacheUpdateRequest) {
< /code>
Вот моя конфигурация Kafka, которая загружается при запуске этого приложения через CCM2 < /p>
{
"configResolution": {
"resolved": {
"clusterPropertyName": "test-cluster",
"displayName": "cacheUpdateRequestListenerID",
"groupId": "test-cache-update-event-cg",
"fetchMinBytes": 100000,
"retryMultiplier": 1,
"enabledHeaderRegex": ".*",
"heartbeatIntervalInMillis": 3000,
"valueDeserializer": "org.apache.kafka.common.serialization.StringDeserializer",
"maxPartitionFetchBytes": 2097152,
"connectionMaxIdleTimeInMillis": 540000,
"maxPollIntervalInMillis": 1200000,
"clientId": "cacheUpdateClient",
"requestTimeoutIntervalInMillis": 30000,
"flowName": "CACHE_UPDATE_ASYNC",
"autoStart": true,
"concurrency": 1,
"isEnableAutoCommit": true,
"consumerAutoOffsetReset": "latest",
"fetchMaxBytes": 5000000,
"retryMaxAttempts": 1,
"maxPollRecords": 200,
"ackMode": "RECORD",
"topic": "test-cache-update-event-topic",
"keyDeserializer": "org.apache.kafka.common.serialization.StringDeserializer",
"additionalProperties": "{}",
"sessionTimeoutInMillis": 10000,
"nonRetryAbleExceptions": [
""
],
"referenceCircuitBreakers": [
""
],
"retryAbleExceptions": [
""
]
}
}
}
Подробнее здесь: https://stackoverflow.com/questions/797 ... er-in-java