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.RecordDeserializationException;
import org.apache.kafka.common.errors.WakeupException;
import reactor.core.publisher.Mono;
public class ChatRoomChannel implements Runnable
{
private final String topic;
- private final Consumer<String, ChatRoomTo> consumer;
- private final Producer<String, ChatRoomTo> producer;
- private final ChatRoomFactory chatRoomFactory;
+ private final Producer<Integer, CreateChatRoomRequestTo> producer;
+ private final Consumer<Integer, CreateChatRoomRequestTo> consumer;
+ private final ShardingStrategy shardingStrategy;
+ private final ChatMessageChannel chatMessageChannel;
+ private final Clock clock;
+ private final int bufferSize;
private boolean running;
{
try
{
- ConsumerRecords<String, ChatRoomTo> records = consumer.poll(Duration.ofMinutes(5));
+ ConsumerRecords<Integer, CreateChatRoomRequestTo> records = consumer.poll(Duration.ofMinutes(5));
log.info("Fetched {} messages", records.count());
- for (ConsumerRecord<String, ChatRoomTo> record : records)
+ for (ConsumerRecord<Integer, CreateChatRoomRequestTo> record : records)
{
- UUID id = record.value().getId();
- String name = record.value().getName();
- chatRoomFactory.createChatRoom(id, name);
+ createChatRoom(record.value().toChatRoomInfo());
}
}
catch (WakeupException e)
{
- }
- catch (RecordDeserializationException e)
- {
+ log.info("Received WakeupException, exiting!");
+ running = false;
}
}
+
+ log.info("Exiting normally");
}
- Mono<ChatRoomInfo> sendCreateChatRoomRequest(
- UUID chatRoomId,
- String name)
- {
- ChatRoomTo chatRoomTo = ChatRoomTo.of(chatRoomId, name);
- return Mono.create(sink ->
- {
- ProducerRecord<String, ChatRoomTo> record =
- new ProducerRecord<>(
- topic,
- chatRoomId.toString(),
- chatRoomTo);
- producer.send(record, ((metadata, exception) ->
- {
- if (metadata != null)
- {
- log.info("Successfully send chreate-request for chat room: {}", chatRoomTo);
- sink.success(chatRoomTo.toChatRoomInfo());
- }
- else
- {
- // On send-failure
- log.error(
- "Could not send create-request for chat room (id={}, name={}): {}",
- chatRoomId,
- name,
- exception);
- sink.error(exception);
- }
- }));
- });
+ void createChatRoom(ChatRoomInfo chatRoomInfo)
+ {
+ UUID id = chatRoomInfo.getId();
+ String name = chatRoomInfo.getName();
+ int shard = chatRoomInfo.getShard();
+ log.info("Creating ChatRoom {} with buffer-size {}", id, bufferSize);
+ KafkaChatRoomService service = new KafkaChatRoomService(chatMessageChannel, id);
+ ChatRoom chatRoom = new ChatRoom(id, name, shard, clock, service, bufferSize);
+ chatMessageChannel.putChatRoom(chatRoom);
}
}