SAMZA-2297: InMemorySystemAdmin offsets are off-by-one in some cases - #1133
Conversation
mynameborat
left a comment
There was a problem hiding this comment.
Thank you for the fix and unit tests :)
| String newestOffset = String.valueOf(entry.getValue().size()); | ||
| String upcomingOffset = String.valueOf(entry.getValue().size() + 1); | ||
| List<IncomingMessageEnvelope> messages = entry.getValue(); | ||
| String oldestOffset = messages.isEmpty() ? null : "0"; |
There was a problem hiding this comment.
@cameronlee314 Any objection to defaulting this to "0" instead of null in case the SSP is empty? This is consistent with the behavior for Kafka as well where an empty SSP returns (0, null, 0) as (oldest, newest, upcoming) offsets. Will also remove the need to handle nulls in Consumer#register.
There was a problem hiding this comment.
SystemStreamMetadata.SystemStreamPartitionMetadata.getOldestOffset is documented with "A null value means the stream is empty". That's why I used null here. Does that mean that Kafka is not following the API?
There was a problem hiding this comment.
Seems like it. https://github.com/apache/samza/blob/master/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java#L403
I need to rely on this information during changelog restore for transactional state. It doesn't make sense to change the behavior for Kafka as part of the transactional state changes. Since InMemoryStore is used as a changelog for tests, I'll relax the assertion/validation for non-null starting offsets in the restore path. Thanks for the pointer.
Added unit tests.
Ran a local build to make sure existing usages are still working.
Changed
TestSamzaSqlEndToEnd.testEndToEndStreamTableRightJoin(which uses intermediate streams) to use in-memory system. Failed before this change, succeeded after this change.