Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

Table of Contents

Status

Current state:  Under DiscussionAccepted

Discussion thread: TODO here 

Voting thread:here

JIRA: TODO KAFKA-20229 

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

...

We have the RemoteClusterUtils API to enable users to perform some operations of their MirrorMaker environmentenvironments. One method from this API, translateOffsets(), allows to manually translate consumer group offsets. This is useful when MirrorCheckpointConnector is not configured to automatically sync offsets or when dealing with a disaster and unplanned failover. This method works by fully reading the checkpoints topic which is populated by MirrorCheckpointConnector, so this can be an expensive call and it can take a long time depending on the size of the topic. This method takes a single consumer group id as input, so it makes it unsuitable in cases when the offsets of several consumer groups have to be translated as each call re-reads the full topic.

Since the full topic is read anyway, we should have a method to translate the committed offsets of multiple consumer groups at the same time.

Public Interfaces

A new translateOffsetsnew translateOffsets() method in RemoteClusterUtils , that takes a set of consumer group Idsregex pattern to specify the consumer groups:

Code Block
languagejava
/**
 * Translates remote consumer groups' offsets into corresponding local offsets. Topics are automatically
 *  renamed according to the configured {@link ReplicationPolicy}.
 *  @param properties Map of properties to instantiate a {@link MirrorClient}
 *  @param remoteClusterAlias The alias of the remote cluster
 *  @param consumerGroupIdsconsumerGroupPattern The regex pattern setspecifying ofthe consumer group Ids
  groups to translate offsets for
 *  @param timeout The maximum time to block when consuming from the checkpoints topic
 *  @throws IllegalArgumentException If any of the arguments are null
 */
public static Map<String, Map<TopicPartition, OffsetAndMetadata>> translateOffsets(Map<String, Object> properties, String remoteClusterAlias, Set<String>Pattern consumerGroupIdsconsumerGroupPattern, Duration timeout) {
}

The matching method in MirrorClient:

Code Block
languagejava
/**
 * Translates remote consumer groups' offsets into corresponding local offsets. Topics are automatically
 * renamed according to the ReplicationPolicy.
 * @param consumerGroupIdsconsumerGroupPattern The regex setpattern specifying ofthe consumer group Ids groups to translate offsets for
 * @param remoteClusterAlias The alias of remote cluster
 * @param timeout The maximum time to block when consuming from the checkpoints topic
 * @throws IllegalArgumentException If any of the arguments are null
 */
public Map<String, Map<TopicPartition, OffsetAndMetadata>> remoteConsumerOffsets(Set<String>Pattern consumerGroupIdsconsumerGroupPattern, String remoteClusterAlias, Duration timeout) {
}

...

The existing MirrorClient.remoteConsumerOffsets() will invoke the new method with a Set containing just a single regex just matching the specified consumer group Id.  

Compatibility, Deprecation, and Migration Plan

  • This is adding new methods. There are no behavior changes to existing methods or logic.

Test Plan

The new methods will be tested using unit tests. This will require some refactoring in MirrorClient to make the remoteConsumerOffsets() method testable (we currently don't have tests for that method).

...

  • Deprecate the existing RemoteClusterUtils.translateOffsets() and MirrorClient.remoteConsumerOffsets() methods: Internally these will use the same logic as the new methods, so we can keep them.
  • Specify the desired consumer groups via a collection: This limits the ability of translating all consumer group offsets.