X-Git-Url: https://juplo.de/gitweb/?a=blobdiff_plain;f=src%2Fmain%2Fjava%2Fde%2Fjuplo%2Fkafka%2Fchat%2Fbackend%2Fimplementation%2Fkafka%2FInfoChannel.java;h=3fe15c44a61410f17a8b2723c92e77fb2e104cce;hb=284e5c04e30bcfecfe2bca1252d247dbc6baaecd;hp=cf4d03c280225405aafe47ff1fd2b023454a6c76;hpb=a12ced5d2e733b957b106e8393f115941d983823;p=demos%2Fkafka%2Fchat diff --git a/src/main/java/de/juplo/kafka/chat/backend/implementation/kafka/InfoChannel.java b/src/main/java/de/juplo/kafka/chat/backend/implementation/kafka/InfoChannel.java index cf4d03c2..3fe15c44 100644 --- a/src/main/java/de/juplo/kafka/chat/backend/implementation/kafka/InfoChannel.java +++ b/src/main/java/de/juplo/kafka/chat/backend/implementation/kafka/InfoChannel.java @@ -68,7 +68,7 @@ public class InfoChannel implements Runnable } - boolean loadInProgress() + boolean isLoadInProgress() { return IntStream .range(0, numShards) @@ -196,7 +196,11 @@ public class InfoChannel implements Runnable { ConsumerRecords records = consumer.poll(Duration.ofMinutes(1)); log.debug("Fetched {} messages", records.count()); - handleMessages(records); + for (ConsumerRecord record : records) + { + handleMessage(record); + updateNextOffset(record.partition(), record.offset() + 1); + } } catch (WakeupException e) { @@ -208,47 +212,47 @@ public class InfoChannel implements Runnable log.info("Exiting normally"); } - private void handleMessages(ConsumerRecords records) + private void updateNextOffset(int partition, long nextOffset) { - for (ConsumerRecord record : records) - { - switch (record.value().getType()) - { - case EVENT_CHATROOM_CREATED: - EventChatRoomCreated eventChatRoomCreated = - (EventChatRoomCreated) record.value(); - createChatRoom(eventChatRoomCreated.toChatRoomInfo()); - break; - - case EVENT_SHARD_ASSIGNED: - EventShardAssigned eventShardAssigned = - (EventShardAssigned) record.value(); - log.info( - "Shard {} was assigned to {}", - eventShardAssigned.getShard(), - eventShardAssigned.getUri()); - shardOwners[eventShardAssigned.getShard()] = eventShardAssigned.getUri(); - break; - - case EVENT_SHARD_REVOKED: - EventShardRevoked eventShardRevoked = - (EventShardRevoked) record.value(); - log.info( - "Shard {} was revoked from {}", - eventShardRevoked.getShard(), - eventShardRevoked.getUri()); - shardOwners[eventShardRevoked.getShard()] = null; - break; - - default: - log.debug( - "Ignoring message for key={} with offset={}: {}", - record.key(), - record.offset(), - record.value()); - } + this.nextOffset[partition] = nextOffset; + } - nextOffset[record.partition()] = record.offset() + 1; + private void handleMessage(ConsumerRecord record) + { + switch (record.value().getType()) + { + case EVENT_CHATROOM_CREATED: + EventChatRoomCreated eventChatRoomCreated = + (EventChatRoomCreated) record.value(); + createChatRoom(eventChatRoomCreated.toChatRoomInfo()); + break; + + case EVENT_SHARD_ASSIGNED: + EventShardAssigned eventShardAssigned = + (EventShardAssigned) record.value(); + log.info( + "Shard {} was assigned to {}", + eventShardAssigned.getShard(), + eventShardAssigned.getUri()); + shardOwners[eventShardAssigned.getShard()] = eventShardAssigned.getUri(); + break; + + case EVENT_SHARD_REVOKED: + EventShardRevoked eventShardRevoked = + (EventShardRevoked) record.value(); + log.info( + "Shard {} was revoked from {}", + eventShardRevoked.getShard(), + eventShardRevoked.getUri()); + shardOwners[eventShardRevoked.getShard()] = null; + break; + + default: + log.debug( + "Ignoring message for key={} with offset={}: {}", + record.key(), + record.offset(), + record.value()); } }