MessagesProcessing.java 2.6 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182
  1. package com.provectus.kafka.ui.emitter;
  2. import com.provectus.kafka.ui.model.TopicMessageDTO;
  3. import com.provectus.kafka.ui.model.TopicMessageEventDTO;
  4. import com.provectus.kafka.ui.model.TopicMessagePhaseDTO;
  5. import com.provectus.kafka.ui.serdes.ConsumerRecordDeserializer;
  6. import java.util.function.Predicate;
  7. import javax.annotation.Nullable;
  8. import lombok.extern.slf4j.Slf4j;
  9. import org.apache.kafka.clients.consumer.ConsumerRecord;
  10. import org.apache.kafka.clients.consumer.ConsumerRecords;
  11. import org.apache.kafka.common.utils.Bytes;
  12. import reactor.core.publisher.FluxSink;
  13. @Slf4j
  14. public class MessagesProcessing {
  15. private final ConsumingStats consumingStats = new ConsumingStats();
  16. private long sentMessages = 0;
  17. private int filterApplyErrors = 0;
  18. private final ConsumerRecordDeserializer deserializer;
  19. private final Predicate<TopicMessageDTO> filter;
  20. private final @Nullable Integer limit;
  21. public MessagesProcessing(ConsumerRecordDeserializer deserializer,
  22. Predicate<TopicMessageDTO> filter,
  23. @Nullable Integer limit) {
  24. this.deserializer = deserializer;
  25. this.filter = filter;
  26. this.limit = limit;
  27. }
  28. boolean limitReached() {
  29. return limit != null && sentMessages >= limit;
  30. }
  31. void sendMsg(FluxSink<TopicMessageEventDTO> sink, ConsumerRecord<Bytes, Bytes> rec) {
  32. if (!sink.isCancelled() && !limitReached()) {
  33. TopicMessageDTO topicMessage = deserializer.deserialize(rec);
  34. try {
  35. if (filter.test(topicMessage)) {
  36. sink.next(
  37. new TopicMessageEventDTO()
  38. .type(TopicMessageEventDTO.TypeEnum.MESSAGE)
  39. .message(topicMessage)
  40. );
  41. sentMessages++;
  42. }
  43. } catch (Exception e) {
  44. filterApplyErrors++;
  45. log.trace("Error applying filter for message {}", topicMessage);
  46. }
  47. }
  48. }
  49. int sentConsumingInfo(FluxSink<TopicMessageEventDTO> sink,
  50. ConsumerRecords<Bytes, Bytes> polledRecords,
  51. long elapsed) {
  52. if (!sink.isCancelled()) {
  53. return consumingStats.sendConsumingEvt(sink, polledRecords, elapsed, filterApplyErrors);
  54. }
  55. return 0;
  56. }
  57. void sendFinishEvent(FluxSink<TopicMessageEventDTO> sink) {
  58. if (!sink.isCancelled()) {
  59. consumingStats.sendFinishEvent(sink, filterApplyErrors);
  60. }
  61. }
  62. void sendPhase(FluxSink<TopicMessageEventDTO> sink, String name) {
  63. if (!sink.isCancelled()) {
  64. sink.next(
  65. new TopicMessageEventDTO()
  66. .type(TopicMessageEventDTO.TypeEnum.PHASE)
  67. .phase(new TopicMessagePhaseDTO().name(name))
  68. );
  69. }
  70. }
  71. }