Apache Kafka: large retention time vs. fast read of last value

Viewed 209

Dear Apache Kafka friends,

I have a use case for which I am looking for an elegant solution:

Data is published in a Kafka-Topic at a relatively high rate. There are two competing requirements

  • all records should be kept for 7 days (which is configured by min.compaction.lag)
  • applications should read the "last status" from the topic during their initialization phase

LogCompaction is enabled in order for the "last state" to be available in the topic. Now comes the problem. If an application wants to initialize itself from the topic, it has to read a lot of records to get the last state for all keys (the entire topic content must be processed). But this is not performant possible with the amount of records.

Idea

A streaming process streams the data of the topic into a corresponding ShortTerm topic which has a much shorter min.compaction.lag time (1 hour). The applications initialize themselves from this topic.

Risk

The streaming process is a potential source of errors. If it temporarily fails, the applications will no longer receive the latest status.

My Question

Are there any other possible solutions to satisfy the two requirements. Did I maybe miss a Kafa concept that helps to handle these competing requirements?

Any contribution is welcome. Thank you all.

2 Answers

If you don't have a strict guarantee how frequently each key will be updated, you cannot do anything else as you proposed.

To avoid the risk that the downstream app does not get new updates (because the data replication jobs stalls), I would recommend to only bootstrap an app from the short term topic, and let it consume from the original topic afterwards. To not miss any updates, you can sync the switch over as follows:

  1. On app startup, get the replication job's committed offsets from the original topic.
  2. Get the short term topic's current end-offsets (because the replication job will continue to write data, you just need a fixed stopping point).
  3. Consume the short term topic from beginning to the captured end offsets.
  4. Resume consuming from the original topic using the captured committed offsets (from step 1) as start point.

This way, you might read some messages twice, but you won't lose any updates.

To me, the two requirements you have mentioned together with the requirement for new consumers are not competing. In fact, I do not see any reason why you should keep a message of an outdated key in your topic for 7 days, because

  • New consumers are only interested in the latest message of a key.
  • Already existing consumers will have processed the message within 1 hour (as taken from your comments).

Therefore, my understanding is that your requirement "all records should be kept for 7 days" can be replaced by "each consumer should have enough time to consume the message & the latest message for each key should be kept for 7 days".

Please correct me if I am wrong and explain which consumer actually does need "all records for 7 days".

If that is the case you could do the following:

  1. Enable log compaction as well as time-based retention to 7 days for this topic
  2. Fine-tune the compaction frequency to be very eager, meaning to keep as little as possible outdated messages for a key.
  3. Set min.compaction.lag to 1 hour such that all consumers have the chance to keep up.

That way, new consumers will read (almost) only the latest message for each key. If that is not performant enough, you can try increasing the partitions and consumer threads of your consumer groups.

Related