|
1 | 1 | package org.folio.service.event; |
2 | 2 |
|
| 3 | +import static io.vertx.core.Future.succeededFuture; |
3 | 4 | import static org.apache.logging.log4j.LogManager.getLogger; |
4 | | -import static org.folio.service.event.EntityChangedEventPublisherFactory.requestEventPublisher; |
5 | 5 |
|
6 | 6 | import java.util.Map; |
7 | 7 |
|
|
11 | 11 | import org.folio.kafka.SimpleKafkaProducerManager; |
12 | 12 | import org.folio.kafka.services.KafkaEnvironmentProperties; |
13 | 13 | import org.folio.kafka.services.KafkaProducerRecordBuilder; |
14 | | -import org.folio.rest.jaxrs.model.Request; |
15 | 14 | import org.folio.rest.tools.utils.TenantTool; |
16 | 15 |
|
17 | 16 | import io.vertx.core.Context; |
@@ -42,27 +41,51 @@ public Future<Void> publish(K key, DomainEvent<T> event, Map<String, String> oka |
42 | 41 | log.info("publish:: key = {}, eventId = {}, type = {}, topic = {}", key, event.getId(), |
43 | 42 | event.getType(), kafkaTopic); |
44 | 43 |
|
45 | | - KafkaProducerRecord<K, String> producerRecord = |
46 | | - new KafkaProducerRecordBuilder<K, DomainEvent<T>>(TenantTool.tenantId(okapiHeaders)) |
47 | | - .key(key) |
48 | | - .value(event) |
49 | | - .topic(kafkaTopic) |
50 | | - .propagateOkapiHeaders(okapiHeaders) |
51 | | - .build(); |
| 44 | + KafkaProducerRecord<K, String> producerRecord = buildProducerRecord(key, event, okapiHeaders); |
52 | 45 | log.info("publish:: kafkaRecord = [{}]", producerRecord); |
53 | 46 |
|
54 | | - KafkaProducer<K, String> producer = getOrCreateProducer(); |
55 | | - log.info("publish:: Producer created, sending the record..."); |
| 47 | + KafkaProducer<K, String> producer = null; |
| 48 | + try { |
| 49 | + producer = getOrCreateProducer(); |
| 50 | + log.info("publish:: Producer created, sending the record..."); |
| 51 | + send(producer, key, producerRecord); |
| 52 | + } catch (Exception e) { |
| 53 | + log.error("publish:: Failed to initiate send for domain event with key [{}], kafka record [{}]", |
| 54 | + key, producerRecord, e); |
| 55 | + if (producer != null) { |
| 56 | + log.info("publish:: Producer is not null, trying to close. Event key: {}.", key); |
| 57 | + producer.close(); |
| 58 | + } |
| 59 | + failureHandler.handle(e, producerRecord); |
| 60 | + } |
| 61 | + |
| 62 | + return succeededFuture(); |
| 63 | + } |
| 64 | + |
| 65 | + private KafkaProducerRecord<K, String> buildProducerRecord(K key, DomainEvent<T> event, |
| 66 | + Map<String, String> okapiHeaders) { |
| 67 | + |
| 68 | + return new KafkaProducerRecordBuilder<K, DomainEvent<T>>(TenantTool.tenantId(okapiHeaders)) |
| 69 | + .key(key) |
| 70 | + .value(event) |
| 71 | + .topic(kafkaTopic) |
| 72 | + .propagateOkapiHeaders(okapiHeaders) |
| 73 | + .build(); |
| 74 | + } |
| 75 | + |
| 76 | + private void send(KafkaProducer<K, String> producer, K key, |
| 77 | + KafkaProducerRecord<K, String> producerRecord) { |
56 | 78 |
|
57 | | - return producer.send(producerRecord) |
58 | | - .onSuccess(r -> log.info("publish:: Succeeded sending domain event with key [{}], " + |
| 79 | + producer.send(producerRecord) |
| 80 | + .onSuccess(r -> log.info("send:: Succeeded sending domain event with key [{}], " + |
59 | 81 | "kafka record [{}]", key, producerRecord)) |
60 | | - .<Void>mapEmpty() |
61 | 82 | .onFailure(cause -> { |
62 | | - log.error("publish:: Unable to send domain event with key [{}], kafka record [{}]", |
| 83 | + log.error("send:: Unable to send domain event with key [{}], kafka record [{}]", |
63 | 84 | key, producerRecord, cause); |
64 | 85 | failureHandler.handle(cause, producerRecord); |
65 | | - }); |
| 86 | + }) |
| 87 | + .eventually(() -> producer.flush()) |
| 88 | + .eventually(() -> producer.close()); |
66 | 89 | } |
67 | 90 |
|
68 | 91 | private KafkaProducer<K, String> getOrCreateProducer() { |
|
0 commit comments