projects
/
demos
/
kafka
/
chat
/ blobdiff
commit
grep
author
committer
pickaxe
?
search:
re
summary
|
shortlog
|
log
|
commit
|
commitdiff
|
tree
raw
|
inline
| side by side
NG
[demos/kafka/chat]
/
src
/
main
/
java
/
de
/
juplo
/
kafka
/
chat
/
backend
/
persistence
/
kafka
/
ChatRoomChannel.java
diff --git
a/src/main/java/de/juplo/kafka/chat/backend/persistence/kafka/ChatRoomChannel.java
b/src/main/java/de/juplo/kafka/chat/backend/persistence/kafka/ChatRoomChannel.java
index
9ea23b1
..
57deaf2
100644
(file)
--- a/
src/main/java/de/juplo/kafka/chat/backend/persistence/kafka/ChatRoomChannel.java
+++ b/
src/main/java/de/juplo/kafka/chat/backend/persistence/kafka/ChatRoomChannel.java
@@
-14,18
+14,16
@@
import reactor.core.publisher.Mono;
import java.time.*;
import java.util.List;
import java.time.*;
import java.util.List;
-import java.util.Optional;
import java.util.UUID;
import java.util.UUID;
-import java.util.concurrent.Callable;
@RequiredArgsConstructor
@Slf4j
@RequiredArgsConstructor
@Slf4j
-public class ChatRoomChannel implements
Callable<Optional<Exception>>
+public class ChatRoomChannel implements
Runnable
{
private final String topic;
{
private final String topic;
- private final Producer<Integer, C
hatRoom
To> producer;
- private final Consumer<Integer, C
hatRoom
To> consumer;
+ private final Producer<Integer, C
reateChatRoomRequest
To> producer;
+ private final Consumer<Integer, C
reateChatRoomRequest
To> consumer;
private final ShardingStrategy shardingStrategy;
private final ChatMessageChannel chatMessageChannel;
private final Clock clock;
private final ShardingStrategy shardingStrategy;
private final ChatMessageChannel chatMessageChannel;
private final Clock clock;
@@
-39,21
+37,21
@@
public class ChatRoomChannel implements Callable<Optional<Exception>>
String name)
{
int shard = this.shardingStrategy.selectShard(chatRoomId);
String name)
{
int shard = this.shardingStrategy.selectShard(chatRoomId);
- C
hatRoomTo chatRoomTo = ChatRoomTo.of(chatRoomId
, name, shard);
+ C
reateChatRoomRequestTo createChatRoomRequestTo = CreateChatRoomRequestTo.of(chatRoomId.toString()
, name, shard);
return Mono.create(sink ->
{
return Mono.create(sink ->
{
- ProducerRecord<Integer, C
hatRoom
To> record =
+ ProducerRecord<Integer, C
reateChatRoomRequest
To> record =
new ProducerRecord<>(
topic,
shard,
new ProducerRecord<>(
topic,
shard,
- c
hatRoom
To);
+ c
reateChatRoomRequest
To);
producer.send(record, ((metadata, exception) ->
{
if (metadata != null)
{
producer.send(record, ((metadata, exception) ->
{
if (metadata != null)
{
- log.info("Successfully send chreate-request for chat room: {}", c
hatRoom
To);
- sink.success(c
hatRoom
To.toChatRoomInfo());
+ log.info("Successfully send chreate-request for chat room: {}", c
reateChatRoomRequest
To);
+ sink.success(c
reateChatRoomRequest
To.toChatRoomInfo());
}
else
{
}
else
{
@@
-70,7
+68,7
@@
public class ChatRoomChannel implements Callable<Optional<Exception>>
}
@Override
}
@Override
- public
Optional<Exception> call
()
+ public
void run
()
{
consumer.assign(List.of(new TopicPartition(topic, 0)));
{
consumer.assign(List.of(new TopicPartition(topic, 0)));
@@
-80,10
+78,10
@@
public class ChatRoomChannel implements Callable<Optional<Exception>>
{
try
{
{
try
{
- ConsumerRecords<Integer, C
hatRoom
To> records = consumer.poll(Duration.ofMinutes(5));
+ ConsumerRecords<Integer, C
reateChatRoomRequest
To> records = consumer.poll(Duration.ofMinutes(5));
log.info("Fetched {} messages", records.count());
log.info("Fetched {} messages", records.count());
- for (ConsumerRecord<Integer, C
hatRoom
To> record : records)
+ for (ConsumerRecord<Integer, C
reateChatRoomRequest
To> record : records)
{
createChatRoom(record.value().toChatRoomInfo());
}
{
createChatRoom(record.value().toChatRoomInfo());
}
@@
-93,15
+91,9
@@
public class ChatRoomChannel implements Callable<Optional<Exception>>
log.info("Received WakeupException, exiting!");
running = false;
}
log.info("Received WakeupException, exiting!");
running = false;
}
- catch (Exception e)
- {
- log.error("Exiting abnormally!");
- return Optional.of(e);
- }
}
log.info("Exiting normally");
}
log.info("Exiting normally");
- return Optional.empty();
}
}