+++ /dev/null
-package de.juplo.kafka.chat.backend.persistence.kafka;
-
-import de.juplo.kafka.chat.backend.domain.*;
-import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.kafka.clients.consumer.Consumer;
-import org.apache.kafka.clients.consumer.ConsumerRecord;
-import org.apache.kafka.clients.consumer.ConsumerRecords;
-import org.apache.kafka.clients.producer.Producer;
-import org.apache.kafka.clients.producer.ProducerRecord;
-import org.apache.kafka.common.TopicPartition;
-import org.apache.kafka.common.errors.WakeupException;
-import reactor.core.publisher.Flux;
-import reactor.core.publisher.Mono;
-
-import java.time.*;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.UUID;
-import java.util.stream.IntStream;
-
-
-@RequiredArgsConstructor
-@Slf4j
-public class ChatRoomChannel implements Runnable
-{
- private final String topic;
- private final Consumer<String, AbstractTo> consumer;
- private final Map<UUID, ChatRoom> chatrooms = new HashMap<>();
-
- private boolean running;
-
-
- @Override
- public void run()
- {
- consumer.assign(List.of(new TopicPartition(topic, 0)));
-
- running = true;
-
- while (running)
- {
- try
- {
- ConsumerRecords<String, AbstractTo> records = consumer.poll(Duration.ofMinutes(5));
- log.info("Fetched {} messages", records.count());
-
- for (ConsumerRecord<String, AbstractTo> record : records)
- {
- switch (record.value().getType())
- {
- case CHATROOM_INFO:
- createChatRoom((ChatRoomInfoTo) record.value());
- break;
-
- default:
- log.debug(
- "Ignoring message for key {} with offset {}: {}",
- record.key(),
- record.offset(),
- record.value());
- }
- }
- }
- catch (WakeupException e)
- {
- log.info("Received WakeupException, exiting!");
- running = false;
- }
- }
-
- log.info("Exiting normally");
- }
-
-
- void createChatRoom(ChatRoomInfoTo chatRoomInfoTo)
- {
- ChatRoomInfo chatRoomInfo = chatRoomInfoTo.toChatRoomInfo();
- chatrooms.put(chatRoomInfo.getId(), chatRoomInfo);
- }
-
- Flux<ChatRoom> getChatRooms()
- {
- return Flux.fromIterable(chatrooms.values());
- }
-}
@Autowired
ConfigurableApplicationContext context;
- @Autowired
- ChatRoomChannel chatRoomChannel;
@Autowired
Consumer<Integer, CreateChatRoomRequestTo> chatRoomChannelConsumer;
@Autowired
@Override
public void run(ApplicationArguments args) throws Exception
{
- log.info("Starting the consumer for the ChatRoomChannel");
- chatRoomChannelConsumerJob = taskExecutor
- .submitCompletable(chatRoomChannel)
- .exceptionally(e ->
- {
- log.error("The consumer for the ChatRoomChannel exited abnormally!", e);
- return null;
- });
log.info("Starting the consumer for the ChatMessageChannel");
chatMessageChannelConsumerJob = taskExecutor
.submitCompletable(chatMessageChannel)