How to determine Kafka topic partition offset if we haven't consumed any messages yet

Viewed 492

librdkafka contains the function rd_kafka_position which fetches the current offsets for the given topic-partitions. But the comment says:

The \p offset field of each requested partition will be set to the offset
of the last consumed message + 1, or RD_KAFKA_OFFSET_INVALID in case there was
no previous message.

In other words, it won't give you any useful information if no messages have been consumed yet.

I'm interested in the case where I've just subscribed to a topic, and I've already called rd_kafka_seek to either:

  1. seek to a known position (in the case of error recovery), or
  2. seek to the very end of the partition.

What I'd like to know, in this context, is what the offset would be for the next message if one were to be consumed. In other words, in the first case, it should be the same offset that was passed to rd_kafka_seek, and in the second case, it should be 1 plus the offset of the last message that was in the partition when rd_kafka_seek was called.

Unfortunately, just like the comment says, rd_kafka_position doesn't return this information. If no messages have been consumed yet, it gives -1001 (RD_KAFKA_OFFSET_INVALID). If I consume a message and then call rd_kafka_position, it gives the correct offset.

Is there some other function that I can call in order to get the offset before consuming any messages?

2 Answers

I'm not sure what you are after.... "offset" is something that is consumer-specific, in most cases (except for two cases I mention below). It tracks read progress of each specific consumer for each topic/partition, and if there was no reading done by that consumer yet - there is no consumer-specific offset for that topic/partition yet. So, asking for an offset of this consumer in this case does not make any sense - the consumer has not yet read anything, so there is no offset associated with it and it could be started from any offset you wish it to start.

The two main cases when consumer-unrelated offsets are useful are:

  • when you know which offsets in a topic you want to start processing from based on either time of the messages or some custom error logging/reporting you have in your application
  • or when you want to start from either EARLIEST or LATEST available offsets in the topic

If you know what position in a partition you want a consumer to start reading from - you just seek to that position and let your consumer start consuming messages from then on. And then you can track progress of this consumer by asking what offset it is at at any point in time ....

And if you want to start either from earliest or latest position - you can find out what that position is (using KAfkaAdminClient.listOffsets(), for example, in 2.5.x version - that is in Java, I don't know what is an equivalent method in Python), and then again seek to that position and start you consumer from it.

So, in brief, again, you can only expect to get a correct offset for a consumer if it has read anything from the topic; otherwise - the only consumer-unrelated meaningful information would be the earliest, latest or some specific (known) offsets determined by you

Related