From ce84a646697e2f67cd6f08b9b6fe58d4b5f1d18a Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Tue, 7 Jul 2026 14:58:49 +0530 Subject: [PATCH 1/2] KIP-1266: Send keyed events to the remote log metadata topic - This is a draft proposal to implement KIP-1266. The approach taken is to have limited set of changes and prepare the topic to enable compaction. The actual decision to enable compaction or not in the __remote_log_metdata topic is left to the users. We can provide a script that reads all the record keys and outputs the count of null and non-null key record count, then update the topic config to compact if required. --- .../log/remote/storage/RemoteLogMetadata.java | 4 + .../storage/RemoteLogSegmentMetadata.java | 5 ++ .../storage/RemoteLogSegmentMetadataKey.java | 43 +++++++++++ .../RemoteLogSegmentMetadataUpdate.java | 5 ++ .../remote/metadata/storage/ConsumerTask.java | 18 +++-- .../metadata/storage/ProducerManager.java | 62 +++++++++++++-- .../TopicBasedRemoteLogMetadataManager.java | 2 +- .../RemoteLogSegmentMetadataKeySerde.java | 75 +++++++++++++++++++ 8 files changed, 200 insertions(+), 14 deletions(-) create mode 100644 storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataKey.java create mode 100644 storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadata.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadata.java index b6af987fc946e..e3c98ca226b19 100644 --- a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadata.java +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadata.java @@ -60,4 +60,8 @@ public int brokerId() { * @return TopicIdPartition for which this event is generated. */ public abstract TopicIdPartition topicIdPartition(); + + public RemoteLogSegmentMetadataKey metadataKey() { + return null; + } } diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadata.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadata.java index 33dc229f616bd..5a1f068c56301 100644 --- a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadata.java +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadata.java @@ -325,6 +325,11 @@ public TopicIdPartition topicIdPartition() { return remoteLogSegmentId.topicIdPartition(); } + @Override + public RemoteLogSegmentMetadataKey metadataKey() { + return RemoteLogSegmentMetadataKey.of(remoteLogSegmentId, state); + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataKey.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataKey.java new file mode 100644 index 0000000000000..198696e80c6e4 --- /dev/null +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataKey.java @@ -0,0 +1,43 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.server.log.remote.storage; + +import org.apache.kafka.common.Uuid; + +import java.util.Objects; + +/** + * Structured key for remote log segment metadata Kafka records. + */ +public record RemoteLogSegmentMetadataKey(Uuid topicId, int partition, Uuid segmentId, byte stateId) { + + public RemoteLogSegmentMetadataKey(Uuid topicId, int partition, Uuid segmentId, byte stateId) { + this.topicId = Objects.requireNonNull(topicId, "topicId cannot be null"); + this.partition = partition; + this.segmentId = Objects.requireNonNull(segmentId, "segmentId cannot be null"); + this.stateId = stateId; + } + + public static RemoteLogSegmentMetadataKey of(RemoteLogSegmentId segmentId, RemoteLogSegmentState state) { + return new RemoteLogSegmentMetadataKey( + segmentId.topicIdPartition().topicId(), + segmentId.topicIdPartition().partition(), + segmentId.id(), + state.id() + ); + } +} diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataUpdate.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataUpdate.java index 373b13ebab1b0..68086b22fc092 100644 --- a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataUpdate.java +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogSegmentMetadataUpdate.java @@ -89,6 +89,11 @@ public TopicIdPartition topicIdPartition() { return remoteLogSegmentId.topicIdPartition(); } + @Override + public RemoteLogSegmentMetadataKey metadataKey() { + return RemoteLogSegmentMetadataKey.of(remoteLogSegmentId, state); + } + @Override public boolean equals(Object o) { if (this == o) { diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java index cdaff1f26d181..70d6c9f763418 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ConsumerTask.java @@ -168,13 +168,17 @@ void closeConsumer() { } private void processConsumerRecord(ConsumerRecord record) { - final RemoteLogMetadata remoteLogMetadata = serde.deserialize(record.value()); - if (shouldProcess(remoteLogMetadata, record.offset())) { - remotePartitionMetadataEventHandler.handleRemoteLogMetadata(remoteLogMetadata); - readOffsetsByUserTopicPartition.put(remoteLogMetadata.topicIdPartition(), record.offset()); - } else { - log.trace("The event {} is skipped because it is either already processed or not assigned to this consumer", - remoteLogMetadata); + byte[] value = record.value(); + // skip the tombstone records + if (value != null) { + final RemoteLogMetadata remoteLogMetadata = serde.deserialize(value); + if (shouldProcess(remoteLogMetadata, record.offset())) { + remotePartitionMetadataEventHandler.handleRemoteLogMetadata(remoteLogMetadata); + readOffsetsByUserTopicPartition.put(remoteLogMetadata.topicIdPartition(), record.offset()); + } else { + log.trace("The event {} is skipped because it is either already processed or not assigned to this consumer", + remoteLogMetadata); + } } log.trace("Updating consumed offset: {} for partition {}", record.offset(), record.partition()); readOffsetsByMetadataPartition.put(record.partition(), record.offset()); diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ProducerManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ProducerManager.java index a09518a168020..d5ba89ac75c17 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ProducerManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ProducerManager.java @@ -18,12 +18,18 @@ import org.apache.kafka.clients.producer.Callback; import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.server.log.remote.metadata.storage.serialization.RemoteLogMetadataSerde; +import org.apache.kafka.server.log.remote.metadata.storage.serialization.RemoteLogSegmentMetadataKeySerde; import org.apache.kafka.server.log.remote.storage.RemoteLogMetadata; +import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentId; +import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadataKey; +import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadataUpdate; +import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentState; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,16 +46,28 @@ public class ProducerManager implements Closeable { private static final Logger log = LoggerFactory.getLogger(ProducerManager.class); + private final RemoteLogSegmentMetadataKeySerde keySerde = new RemoteLogSegmentMetadataKeySerde(); private final RemoteLogMetadataSerde serde = new RemoteLogMetadataSerde(); - private final KafkaProducer producer; + private final Producer producer; private final RemoteLogMetadataTopicPartitioner topicPartitioner; private final TopicBasedRemoteLogMetadataManagerConfig rlmmConfig; + private final Callback tombstoneRecordsCallback = (metadata, exception) -> { + if (exception != null) { + log.error("Failed to publish tombstone records for: {}", metadata, exception); + } + }; public ProducerManager(TopicBasedRemoteLogMetadataManagerConfig rlmmConfig, RemoteLogMetadataTopicPartitioner rlmmTopicPartitioner) { + this(rlmmConfig, rlmmTopicPartitioner, new KafkaProducer<>(rlmmConfig.producerProperties())); + } + + public ProducerManager(TopicBasedRemoteLogMetadataManagerConfig rlmmConfig, + RemoteLogMetadataTopicPartitioner rlmmTopicPartitioner, + Producer producer) { this.rlmmConfig = rlmmConfig; - this.producer = new KafkaProducer<>(rlmmConfig.producerProperties()); - topicPartitioner = rlmmTopicPartitioner; + this.topicPartitioner = rlmmTopicPartitioner; + this.producer = producer; } /** @@ -57,7 +75,7 @@ public ProducerManager(TopicBasedRemoteLogMetadataManagerConfig rlmmConfig, * is considered complete. * * @param remoteLogMetadata RemoteLogMetadata to be published - * @return CompletableFuture that completes with RecordMetadata when the message is successfully published + * @return a future with acknowledgement. */ CompletableFuture publishMessage(RemoteLogMetadata remoteLogMetadata) { CompletableFuture future = new CompletableFuture<>(); @@ -80,8 +98,10 @@ CompletableFuture publishMessage(RemoteLogMetadata remoteLogMeta future.complete(metadata); } }; - producer.send(new ProducerRecord<>(rlmmConfig.remoteLogMetadataTopicName(), metadataPartitionNum, null, - serde.serialize(remoteLogMetadata)), callback); + String topic = rlmmConfig.remoteLogMetadataTopicName(); + byte[] serializedKey = keySerde.serializer().serialize(topic, remoteLogMetadata.metadataKey()); + byte[] serializedValue = serde.serialize(remoteLogMetadata); + producer.send(new ProducerRecord<>(topic, metadataPartitionNum, serializedKey, serializedValue), callback); } catch (Exception ex) { future.completeExceptionally(ex); } @@ -89,6 +109,36 @@ CompletableFuture publishMessage(RemoteLogMetadata remoteLogMeta return future; } + /** + * Publishes tombstone records to mark the deletion of remote log segment metadata when the state of the provided + * {@link RemoteLogMetadata} instance is {@link RemoteLogSegmentState#DELETE_SEGMENT_FINISHED}. + * + * @param remoteLogMetadata The {@link RemoteLogMetadata} instance containing metadata information. If the metadata is an instance + * of {@link RemoteLogSegmentMetadataUpdate} and its state is {@link RemoteLogSegmentState#DELETE_SEGMENT_FINISHED}, + * tombstone records are published to indicate its deletion. + */ + public void maybePublishTombstoneRecords(RemoteLogMetadata remoteLogMetadata) { + // TODO: Gate sending tombstone records behind a feature flag. + // Send the tombstone records only when the consumerTask can handle the null values. + if (remoteLogMetadata instanceof RemoteLogSegmentMetadataUpdate metadataUpdate) { + if (metadataUpdate.state() == RemoteLogSegmentState.DELETE_SEGMENT_FINISHED) { + RemoteLogSegmentId remoteLogSegmentId = metadataUpdate.remoteLogSegmentId(); + int metadataPartition = topicPartitioner.metadataPartition(remoteLogSegmentId.topicIdPartition()); + String topic = rlmmConfig.remoteLogMetadataTopicName(); + try { + // Send the tombstone records for all the RemoteLogSegment state to cleanup the expired segment metadata. + for (RemoteLogSegmentState state : RemoteLogSegmentState.values()) { + RemoteLogSegmentMetadataKey metadataKey = RemoteLogSegmentMetadataKey.of(remoteLogSegmentId, state); + byte[] serializedKey = keySerde.serializer().serialize(topic, metadataKey); + producer.send(new ProducerRecord<>(topic, metadataPartition, serializedKey, null), tombstoneRecordsCallback); + } + } catch (Exception ex) { + log.error("Failed to publish tombstone records for: {}", metadataUpdate, ex); + } + } + } + } + public void close() { try { producer.close(Duration.ofSeconds(30)); diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java index 47b8d63b1d368..f3dbea2200b3e 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java @@ -162,7 +162,7 @@ private CompletableFuture storeRemoteLogMetadata(RemoteLogMetadata remoteL } catch (TimeoutException e) { throw new KafkaException(e); } - }); + }).thenAcceptAsync(recordMetadata -> producerManager.maybePublishTombstoneRecords(remoteLogMetadata)); } catch (KafkaException e) { if (e instanceof RetriableException) { throw e; diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java new file mode 100644 index 0000000000000..cef96716a07ec --- /dev/null +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java @@ -0,0 +1,75 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.server.log.remote.metadata.storage.serialization; + +import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.serialization.Deserializer; +import org.apache.kafka.common.serialization.Serde; +import org.apache.kafka.common.serialization.Serializer; +import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadataKey; + +import java.nio.ByteBuffer; + +/** + * Serde for {@link RemoteLogSegmentMetadataKey}. + *

+ * Serializes to a compact 37-byte binary representation: + *

    + *
  • 8 bytes – topicId most-significant bits
  • + *
  • 8 bytes – topicId least-significant bits
  • + *
  • 4 bytes – partition (big-endian int)
  • + *
  • 8 bytes – segmentId most-significant bits
  • + *
  • 8 bytes – segmentId least-significant bits
  • + *
  • 1 byte – stateId
  • + *
+ */ +public class RemoteLogSegmentMetadataKeySerde implements Serde { + + static final int SERIALIZED_SIZE = 8 + 8 + 4 + 8 + 8 + 1; // 37 bytes + + @Override + public Serializer serializer() { + return (topic, key) -> { + if (key == null) return null; + ByteBuffer buf = ByteBuffer.allocate(SERIALIZED_SIZE); + buf.putLong(key.topicId().getMostSignificantBits()); + buf.putLong(key.topicId().getLeastSignificantBits()); + buf.putInt(key.partition()); + buf.putLong(key.segmentId().getMostSignificantBits()); + buf.putLong(key.segmentId().getLeastSignificantBits()); + buf.put(key.stateId()); + return buf.array(); + }; + } + + @Override + public Deserializer deserializer() { + return (topic, data) -> { + if (data == null) return null; + if (data.length != SERIALIZED_SIZE) { + throw new IllegalArgumentException( + "Expected " + SERIALIZED_SIZE + " bytes but got " + data.length); + } + ByteBuffer buf = ByteBuffer.wrap(data); + Uuid topicId = new Uuid(buf.getLong(), buf.getLong()); + int partition = buf.getInt(); + Uuid segmentId = new Uuid(buf.getLong(), buf.getLong()); + byte stateId = buf.get(); + return new RemoteLogSegmentMetadataKey(topicId, partition, segmentId, stateId); + }; + } +} From 7fc40be35806513395e443655be63d5c54d550d3 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Tue, 7 Jul 2026 21:51:36 +0530 Subject: [PATCH 2/2] fix the copilot comments --- .../TopicBasedRemoteLogMetadataManager.java | 2 +- .../RemoteLogSegmentMetadataKeySerde.java | 16 +++++++++++++--- 2 files changed, 14 insertions(+), 4 deletions(-) diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java index f3dbea2200b3e..c8553ab33a38e 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java @@ -162,7 +162,7 @@ private CompletableFuture storeRemoteLogMetadata(RemoteLogMetadata remoteL } catch (TimeoutException e) { throw new KafkaException(e); } - }).thenAcceptAsync(recordMetadata -> producerManager.maybePublishTombstoneRecords(remoteLogMetadata)); + }).thenRunAsync(() -> producerManager.maybePublishTombstoneRecords(remoteLogMetadata)); } catch (KafkaException e) { if (e instanceof RetriableException) { throw e; diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java index cef96716a07ec..e39f5cbbc91db 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/serialization/RemoteLogSegmentMetadataKeySerde.java @@ -27,8 +27,9 @@ /** * Serde for {@link RemoteLogSegmentMetadataKey}. *

- * Serializes to a compact 37-byte binary representation: + * Serializes to a 38-byte binary representation: *

    + *
  • 1 byte – format version (currently {@value #VERSION})
  • *
  • 8 bytes – topicId most-significant bits
  • *
  • 8 bytes – topicId least-significant bits
  • *
  • 4 bytes – partition (big-endian int)
  • @@ -39,13 +40,15 @@ */ public class RemoteLogSegmentMetadataKeySerde implements Serde { - static final int SERIALIZED_SIZE = 8 + 8 + 4 + 8 + 8 + 1; // 37 bytes + static final byte VERSION = 0; + static final int SERIALIZED_SIZE = 1 + 8 + 8 + 4 + 8 + 8 + 1; // 38 bytes @Override public Serializer serializer() { return (topic, key) -> { if (key == null) return null; ByteBuffer buf = ByteBuffer.allocate(SERIALIZED_SIZE); + buf.put(VERSION); buf.putLong(key.topicId().getMostSignificantBits()); buf.putLong(key.topicId().getLeastSignificantBits()); buf.putInt(key.partition()); @@ -59,12 +62,19 @@ public Serializer serializer() { @Override public Deserializer deserializer() { return (topic, data) -> { - if (data == null) return null; + if (data == null) + return null; + if (data.length != SERIALIZED_SIZE) { throw new IllegalArgumentException( "Expected " + SERIALIZED_SIZE + " bytes but got " + data.length); } ByteBuffer buf = ByteBuffer.wrap(data); + byte version = buf.get(); + if (version != VERSION) { + throw new IllegalArgumentException( + "Unsupported RemoteLogSegmentMetadataKey version: " + version); + } Uuid topicId = new Uuid(buf.getLong(), buf.getLong()); int partition = buf.getInt(); Uuid segmentId = new Uuid(buf.getLong(), buf.getLong());