-
Notifications
You must be signed in to change notification settings - Fork 0
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
kafka consumer detail #116
Comments
kafka consumer multi worker method
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(10));
for (ConsumerRecord<String, String> record: records) {
ConsumerWorker worker = new ConsumerWorker(record.value());
executorService.execute(worker);
}
|
Kafka consumer multi thread method
ExcutorService executorSrervice = Executors.newCachedThreadPool();
for (int i=0; i<CONSUMER_COUNT; i++) {
ConsumerWorker worker = new ConsumerWorker(configs, TOPIC_NAME, i);
executorService.execute(worker);
} |
컨슈머 랙
컨슈머 metrics()를 이용한 컨슈머 랙 조회for (Map.Extry<MetricName, ? extends Metric> entry : kafkaConsumer.metrics().entrySet()) {
if ("records-lag-max".equl(etnry.getKey().name()) |
"records-lag".equl(etnry.getKey().name()) |
"records-lag-avg".equl(etnry.getKey().name())) {
Metric metric = entry.getValue();
logger.info("{}:{}", entry.getKey().name(), metric.metricValue());
}
}
카프카 버로우
|
컨슈머 배포 프로세스
무중단 배포
컨슈머 카나리 배포
|
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
multithread consumer
컨슈머 운영 전략
The text was updated successfully, but these errors were encountered: