Skip to content
Merged
Show file tree
Hide file tree
Changes from 8 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -27,9 +27,7 @@
import org.slf4j.LoggerFactory;

import java.io.Closeable;
import java.io.File;
import java.io.IOException;
import java.nio.file.Path;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeoutException;
Expand All @@ -41,8 +39,6 @@
*/
public class ConsumerManager implements Closeable {

public static final String COMMITTED_OFFSETS_FILE_NAME = "_rlmm_committed_offsets";

private static final Logger log = LoggerFactory.getLogger(ConsumerManager.class);
private static final long CONSUME_RECHECK_INTERVAL_MS = 50L;

Expand All @@ -60,15 +56,13 @@ public ConsumerManager(TopicBasedRemoteLogMetadataManagerConfig rlmmConfig,

//Create a task to consume messages and submit the respective events to RemotePartitionMetadataEventHandler.
KafkaConsumer<byte[], byte[]> consumer = new KafkaConsumer<>(rlmmConfig.consumerProperties());
Path committedOffsetsPath = new File(rlmmConfig.logDir(), COMMITTED_OFFSETS_FILE_NAME).toPath();
consumerTask = new ConsumerTask(
consumer,
rlmmConfig.remoteLogMetadataTopicName(),
remotePartitionMetadataEventHandler,
topicPartitioner,
committedOffsetsPath,
time,
60_000L
remotePartitionMetadataEventHandler,
topicPartitioner,
consumer,
100L,
300_000L,
time
);
consumerTaskThread = KafkaThread.nonDaemon("RLMMConsumerTask", consumerTask);
}
Expand Down Expand Up @@ -110,7 +104,7 @@ public void waitTillConsumptionCatchesUp(RecordMetadata recordMetadata,
log.info("Waiting until consumer is caught up with the target partition: [{}]", partition);

// If the current assignment does not have the subscription for this partition then return immediately.
if (!consumerTask.isPartitionAssigned(partition)) {
if (!consumerTask.isMetadataPartitionAssigned(partition)) {
throw new KafkaException("This consumer is not assigned to the target partition " + partition + ". " +
"Partitions currently assigned: " + consumerTask.metadataPartitionsAssigned());
}
Expand All @@ -119,17 +113,17 @@ public void waitTillConsumptionCatchesUp(RecordMetadata recordMetadata,
long startTimeMs = time.milliseconds();
while (true) {
log.debug("Checking if partition [{}] is up to date with offset [{}]", partition, offset);
long receivedOffset = consumerTask.receivedOffsetForPartition(partition).orElse(-1L);
if (receivedOffset >= offset) {
long readOffset = consumerTask.readOffsetForMetadataPartition(partition).orElse(-1L);
if (readOffset >= offset) {
return;
}

log.debug("Expected offset [{}] for partition [{}], but the committed offset: [{}], Sleeping for [{}] to retry again",
offset, partition, receivedOffset, consumeCheckIntervalMs);
offset, partition, readOffset, consumeCheckIntervalMs);

if (time.milliseconds() - startTimeMs > timeoutMs) {
log.warn("Expected offset for partition:[{}] is : [{}], but the committed offset: [{}] ",
partition, receivedOffset, offset);
partition, readOffset, offset);
throw new TimeoutException("Timed out in catching up with the expected offset by consumer.");
}

Expand Down Expand Up @@ -158,7 +152,7 @@ public void removeAssignmentsForPartitions(Set<TopicIdPartition> partitions) {
consumerTask.removeAssignmentsForPartitions(partitions);
}

public Optional<Long> receivedOffsetForPartition(int metadataPartition) {
return consumerTask.receivedOffsetForPartition(metadataPartition);
public Optional<Long> readOffsetForPartition(int metadataPartition) {
return consumerTask.readOffsetForMetadataPartition(metadataPartition);
}
}

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CountDownLatch;

/**
* This class provides an in-memory cache of remote log segment metadata. This maintains the lineage of segments
Expand Down Expand Up @@ -104,6 +105,16 @@ public class RemoteLogMetadataCache {
// https://issues.apache.org/jira/browse/KAFKA-12641
protected final ConcurrentMap<Integer, RemoteLogLeaderEpochState> leaderEpochEntries = new ConcurrentHashMap<>();

private final CountDownLatch initializedLatch = new CountDownLatch(1);

public void markInitialized() {
initializedLatch.countDown();
}

public boolean isInitialized() {
return initializedLatch.getCount() == 0;
}

/**
* Returns {@link RemoteLogSegmentMetadata} if it exists for the given leader-epoch containing the offset and with
* {@link RemoteLogSegmentState#COPY_SEGMENT_FINISHED} state, else returns {@link Optional#empty()}.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,4 +50,9 @@ public abstract void syncLogMetadataSnapshot(TopicIdPartition topicIdPartition,

public abstract void clearTopicPartition(TopicIdPartition topicIdPartition);

public abstract void markInitialized(TopicIdPartition partition);

public abstract boolean isInitialized(TopicIdPartition partition);

public abstract void maybeLoadPartition(TopicIdPartition partition);
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

import org.apache.kafka.common.TopicIdPartition;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.ReplicaNotAvailableException;
import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentId;
import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadata;
import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadataUpdate;
Expand Down Expand Up @@ -151,6 +152,10 @@ private FileBasedRemoteLogMetadataCache getRemoteLogMetadataCache(TopicIdPartiti
throw new RemoteResourceNotFoundException("No resource found for partition: " + topicIdPartition);
}

if (!remoteLogMetadataCache.isInitialized()) {
throw new ReplicaNotAvailableException("Remote log metadata cache is not initialized for partition: " + topicIdPartition);
Comment thread
abhijeetk88 marked this conversation as resolved.
}

return remoteLogMetadataCache;
}

Expand Down Expand Up @@ -180,9 +185,21 @@ public void close() throws IOException {
idToRemoteLogMetadataCache = Collections.emptyMap();
}

@Override
public void maybeLoadPartition(TopicIdPartition partition) {
idToRemoteLogMetadataCache.computeIfAbsent(partition,
topicIdPartition -> new FileBasedRemoteLogMetadataCache(topicIdPartition, partitionLogDirectory(topicIdPartition.topicPartition())));
}
Comment thread
abhijeetk88 marked this conversation as resolved.

@Override
public void markInitialized(TopicIdPartition partition) {
idToRemoteLogMetadataCache.get(partition).markInitialized();
log.trace("Remote log components are initialized for user-partition: {}", partition);
}

@Override
public boolean isInitialized(TopicIdPartition topicIdPartition) {
RemoteLogMetadataCache metadataCache = idToRemoteLogMetadataCache.get(topicIdPartition);
return metadataCache != null && metadataCache.isInitialized();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ public class TopicBasedRemoteLogMetadataManager implements RemoteLogMetadataMana

private RemotePartitionMetadataStore remotePartitionMetadataStore;
private volatile TopicBasedRemoteLogMetadataManagerConfig rlmmConfig;
private volatile RemoteLogMetadataTopicPartitioner rlmmTopicPartitioner;
private volatile RemoteLogMetadataTopicPartitioner rlmTopicPartitioner;
private final Set<TopicIdPartition> pendingAssignPartitions = Collections.synchronizedSet(new HashSet<>());
private volatile boolean initializationFailed;

Expand Down Expand Up @@ -260,12 +260,12 @@ public Iterator<RemoteLogSegmentMetadata> listRemoteLogSegments(TopicIdPartition
}

public int metadataPartition(TopicIdPartition topicIdPartition) {
return rlmmTopicPartitioner.metadataPartition(topicIdPartition);
return rlmTopicPartitioner.metadataPartition(topicIdPartition);
}

// Visible For Testing
public Optional<Long> receivedOffsetForPartition(int metadataPartition) {
return consumerManager.receivedOffsetForPartition(metadataPartition);
public Optional<Long> readOffsetForPartition(int metadataPartition) {
return consumerManager.readOffsetForPartition(metadataPartition);
}

@Override
Expand Down Expand Up @@ -357,7 +357,7 @@ public void configure(Map<String, ?> configs) {
log.info("Started configuring topic-based RLMM with configs: {}", configs);

rlmmConfig = new TopicBasedRemoteLogMetadataManagerConfig(configs);
rlmmTopicPartitioner = new RemoteLogMetadataTopicPartitioner(rlmmConfig.metadataTopicPartitionsCount());
rlmTopicPartitioner = new RemoteLogMetadataTopicPartitioner(rlmmConfig.metadataTopicPartitionsCount());
remotePartitionMetadataStore = new RemotePartitionMetadataStore(new File(rlmmConfig.logDir()).toPath());
configured = true;
log.info("Successfully configured topic-based RLMM with config: {}", rlmmConfig);
Expand Down Expand Up @@ -416,8 +416,8 @@ private void initializeResources() {
// Create producer and consumer managers.
lock.writeLock().lock();
try {
producerManager = new ProducerManager(rlmmConfig, rlmmTopicPartitioner);
consumerManager = new ConsumerManager(rlmmConfig, remotePartitionMetadataStore, rlmmTopicPartitioner, time);
producerManager = new ProducerManager(rlmmConfig, rlmTopicPartitioner);
consumerManager = new ConsumerManager(rlmmConfig, remotePartitionMetadataStore, rlmTopicPartitioner, time);
if (startConsumerThread) {
consumerManager.startConsumerThread();
} else {
Expand Down Expand Up @@ -509,10 +509,8 @@ public TopicBasedRemoteLogMetadataManagerConfig config() {
}

// Visible for testing.
public void startConsumerThread() {
if (consumerManager != null) {
consumerManager.startConsumerThread();
}
void setRlmTopicPartitioner(RemoteLogMetadataTopicPartitioner rlmTopicPartitioner) {
this.rlmTopicPartitioner = Objects.requireNonNull(rlmTopicPartitioner);
}

@Override
Expand Down
Loading