Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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 @@ -5154,6 +5154,8 @@ private static long getOffsetFromSpec(OffsetSpec offsetSpec) {
return ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP;
} else if (offsetSpec instanceof OffsetSpec.LatestTieredSpec) {
return ListOffsetsRequest.LATEST_TIERED_TIMESTAMP;
} else if (offsetSpec instanceof OffsetSpec.EarliestPendingUploadSpec) {
return ListOffsetsRequest.EARLIEST_PENDING_UPLOAD_TIMESTAMP;
}
return ListOffsetsRequest.LATEST_TIMESTAMP;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ public static class LatestSpec extends OffsetSpec { }
public static class MaxTimestampSpec extends OffsetSpec { }
public static class EarliestLocalSpec extends OffsetSpec { }
public static class LatestTieredSpec extends OffsetSpec { }
public static class EarliestPendingUploadSpec extends OffsetSpec { }
public static class TimestampSpec extends OffsetSpec {
private final long timestamp;

Expand Down Expand Up @@ -91,4 +92,13 @@ public static OffsetSpec earliestLocal() {
public static OffsetSpec latestTiered() {
return new LatestTieredSpec();
}

/**
* Used to retrieve the earliest offset of records that are pending upload to remote storage.
* <br/>
* Note: When tiered storage is not enabled, we will return unknown offset.
*/
public static OffsetSpec earliestPendingUpload() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I mean this new public API

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sure, will add.

return new EarliestPendingUploadSpec();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -103,12 +103,17 @@ ListOffsetsRequest.Builder buildBatchedRequest(int brokerId, Set<TopicPartition>
.stream()
.anyMatch(key -> offsetTimestampsByPartition.get(key) == ListOffsetsRequest.LATEST_TIERED_TIMESTAMP);

boolean requireEarliestPendingUploadTimestamp = keys
.stream()
.anyMatch(key -> offsetTimestampsByPartition.get(key) == ListOffsetsRequest.EARLIEST_PENDING_UPLOAD_TIMESTAMP);

int timeoutMs = options.timeoutMs() != null ? options.timeoutMs() : defaultApiTimeoutMs;
return ListOffsetsRequest.Builder.forConsumer(true,
options.isolationLevel(),
supportsMaxTimestamp,
requireEarliestLocalTimestamp,
requireTieredStorageTimestamp)
requireTieredStorageTimestamp,
requireEarliestPendingUploadTimestamp)
.setTargetTimes(new ArrayList<>(topicsByName.values()))
.setTimeoutMs(timeoutMs);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ public class ListOffsetsRequest extends AbstractRequest {

public static final long LATEST_TIERED_TIMESTAMP = -5L;

public static final long EARLIEST_PENDING_UPLOAD_TIMESTAMP = -6L;

public static final int CONSUMER_REPLICA_ID = -1;
public static final int DEBUGGING_REPLICA_ID = -2;

Expand All @@ -58,16 +60,19 @@ public static class Builder extends AbstractRequest.Builder<ListOffsetsRequest>

public static Builder forConsumer(boolean requireTimestamp,
IsolationLevel isolationLevel) {
return forConsumer(requireTimestamp, isolationLevel, false, false, false);
return forConsumer(requireTimestamp, isolationLevel, false, false, false, false);
}

public static Builder forConsumer(boolean requireTimestamp,
IsolationLevel isolationLevel,
boolean requireMaxTimestamp,
boolean requireEarliestLocalTimestamp,
boolean requireTieredStorageTimestamp) {
boolean requireTieredStorageTimestamp,
boolean requireEarliestPendingUploadTimestamp) {
short minVersion = ApiKeys.LIST_OFFSETS.oldestVersion();
if (requireTieredStorageTimestamp)
if (requireEarliestPendingUploadTimestamp)
minVersion = 11;
else if (requireTieredStorageTimestamp)
minVersion = 9;
else if (requireEarliestLocalTimestamp)
minVersion = 8;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,9 @@
// Version 9 enables listing offsets by last tiered offset (KIP-1005).
//
// Version 10 enables async remote list offsets support (KIP-1075)
"validVersions": "1-10",
//
// Version 11 enables listing offsets by earliest pending upload offset (KIP-1023)
"validVersions": "1-11",
"flexibleVersions": "6+",
"latestVersionUnstable": false,
"fields": [
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,9 @@
// Version 9 enables listing offsets by last tiered offset (KIP-1005).
//
// Version 10 enables async remote list offsets support (KIP-1075)
"validVersions": "1-10",
//
// Version 11 enables listing offsets by earliest pending upload offset (KIP-1023)
"validVersions": "1-11",
"flexibleVersions": "6+",
"fields": [
{ "name": "ThrottleTimeMs", "type": "int32", "versions": "2+", "ignorable": true,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8729,6 +8729,34 @@ public void testListOffsetsLatestTierSpecSpecMinVersion() throws Exception {
}
}

@Test
public void testListOffsetsEarliestPendingUploadSpecSpecMinVersion() throws Exception {
Node node = new Node(0, "localhost", 8120);
List<Node> nodes = Collections.singletonList(node);
List<PartitionInfo> pInfos = new ArrayList<>();
pInfos.add(new PartitionInfo("foo", 0, node, new Node[]{node}, new Node[]{node}));
final Cluster cluster = new Cluster(
"mockClusterId",
nodes,
pInfos,
Collections.emptySet(),
Collections.emptySet(),
node);
final TopicPartition tp0 = new TopicPartition("foo", 0);
try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(cluster,
AdminClientConfig.RETRIES_CONFIG, "2")) {

env.kafkaClient().setNodeApiVersions(NodeApiVersions.create());
env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE));

env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.earliestPendingUpload()));

TestUtils.waitForCondition(() -> env.kafkaClient().requests().stream().anyMatch(request ->
request.requestBuilder().apiKey().messageType == ApiMessageType.LIST_OFFSETS && request.requestBuilder().oldestAllowedVersion() == 11
), "no listOffsets request has the expected oldestAllowedVersion");
}
}

private Map<String, FeatureUpdate> makeTestFeatureUpdates() {
return Utils.mkMap(
Utils.mkEntry("test_feature_1", new FeatureUpdate((short) 2, FeatureUpdate.UpgradeType.UPGRADE)),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -127,19 +127,23 @@ public void testListOffsetsRequestOldestVersion() {
.forConsumer(false, IsolationLevel.READ_COMMITTED);

ListOffsetsRequest.Builder maxTimestampRequestBuilder = ListOffsetsRequest.Builder
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, true, false, false);
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, true, false, false, false);

ListOffsetsRequest.Builder requireEarliestLocalTimestampRequestBuilder = ListOffsetsRequest.Builder
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, false, true, false);
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, false, true, false, false);

ListOffsetsRequest.Builder requireTieredStorageTimestampRequestBuilder = ListOffsetsRequest.Builder
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, false, false, true);
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, false, false, true, false);

ListOffsetsRequest.Builder requireEarliestPendingUploadTimestampRequestBuilder = ListOffsetsRequest.Builder
.forConsumer(false, IsolationLevel.READ_UNCOMMITTED, false, false, false, true);

assertEquals((short) 1, consumerRequestBuilder.oldestAllowedVersion());
assertEquals((short) 1, requireTimestampRequestBuilder.oldestAllowedVersion());
assertEquals((short) 2, requestCommittedRequestBuilder.oldestAllowedVersion());
assertEquals((short) 7, maxTimestampRequestBuilder.oldestAllowedVersion());
assertEquals((short) 8, requireEarliestLocalTimestampRequestBuilder.oldestAllowedVersion());
assertEquals((short) 9, requireTieredStorageTimestampRequestBuilder.oldestAllowedVersion());
assertEquals((short) 11, requireEarliestPendingUploadTimestampRequestBuilder.oldestAllowedVersion());
}
}
3 changes: 2 additions & 1 deletion core/src/main/scala/kafka/server/ReplicaManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,8 @@ object ReplicaManager {
ListOffsetsRequest.LATEST_TIMESTAMP -> 1.toShort,
ListOffsetsRequest.MAX_TIMESTAMP -> 7.toShort,
ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP -> 8.toShort,
ListOffsetsRequest.LATEST_TIERED_TIMESTAMP -> 9.toShort
ListOffsetsRequest.LATEST_TIERED_TIMESTAMP -> 9.toShort,
ListOffsetsRequest.EARLIEST_PENDING_UPLOAD_TIMESTAMP -> 11.toShort
)

def createLogReadResult(highWatermark: Long,
Expand Down
Loading