DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
| Table of Contents |
|---|
This page is meant as a template for writing a KIP. To create a KIP choose Tools->Copy on this page and modify with your content and replace the heading with the next KIP number and a description of your issue. Replace anything in italics with your own description.
Status
Current state: Under Discussion
...
We have the RemoteClusterUtils API to enable users to perform some operations of their MirrorMaker environment. One method, translateOffsets(), allows to manually translate consumer 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 and take a long time if the topic is large. This method takes a single consumer group id as input so it makes it unusable 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 offsets the committed offsets of multiple groups at a time.
Public Interfaces
Briefly list any new interfaces that will be introduced as part of this proposal or any existing interfaces that will be removed or changed. The purpose of this section is to concisely call out the public contract that will come along with this feature.
A public interface is any change to the following:
Binary log format
The network protocol and api behavior
Any class in the public packages under clientsConfiguration, especially client configuration
org/apache/kafka/common/serialization
org/apache/kafka/common
org/apache/kafka/common/errors
org/apache/kafka/clients/producer
org/apache/kafka/clients/consumer (eventually, once stable)
Monitoring
Command line tools and arguments
- Anything else that will likely break existing users in some way when they upgrade
Proposed Changes
A new translateOffsets() method in RemoteClusterUtils, that takes a set of consumer group Ids:
| Code Block | ||
|---|---|---|
| ||
/**
* 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 consumerGroupIds The list of consumer group Ids
* @param timeout The maximum time to block when consuming from the checkpoints topic
*/
public static Map<String, Map<TopicPartition, OffsetAndMetadata>> translateOffsets(Map<String, Object> properties, String remoteClusterAlias, Set<String> consumerGroupIds, Duration timeout) {
} |
The matching method in MirrorClient:
| Code Block | ||
|---|---|---|
| ||
/**
* Translates remote consumer groups' offsets into corresponding local offsets. Topics are automatically
* renamed according to the ReplicationPolicy.
* @param consumerGroupIds The remote consumer group Ids
* @param remoteClusterAlias The alias of remote cluster
* @param timeout The maximum time to block when consuming from the checkpoints topic
*/
public Map<String, Map<TopicPartition, OffsetAndMetadata>> remoteConsumerOffsets(Set<String> consumerGroupIds, String remoteClusterAlias, Duration timeout) {
} |
Proposed Changes
The existing MirrorClient.remoteConsumerOffsets() will invoke the new method with a Set containing just a single consumer group Id. Describe the new thing you want to do in appropriate detail. This may be fairly extensive and have large subsections of its own. Or it may be a few sentences. Use judgement based on the scope of the change.
Compatibility, Deprecation, and Migration Plan
- What impact (if any) will there be on existing users?
- If we are changing behavior how will we phase out the older behavior?
- If we need special migration tools, describe them here.
- When will we remove the existing behavior?
Test Plan
Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?
Rejected Alternatives
...
- This is adding new methods. There are no changes to existing methods or logic.
Test Plan
The new methods will be tested using unit tests.
Rejected Alternatives
- 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.