This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 931c67a0f25 Introduce reason codes for segment completion requests
(#19251)
931c67a0f25 is described below
commit 931c67a0f2574568d5d82e1536997eb3f753291e
Author: Arunkumar Ganesan <[email protected]>
AuthorDate: Sat Oct 3 17:02:55 2026 -0400
Introduce reason codes for segment completion requests (#19251)
---
.../protocols/SegmentCompletionProtocol.java | 96 +++++++++++++++++++++-
.../protocols/SegmentCompletionProtocolTest.java | 89 +++++++++++++++++++-
.../resources/LLCSegmentCompletionHandlers.java | 23 ++++--
.../LLCSegmentCompletionHandlersTest.java | 92 +++++++++++++++++++++
.../helix/core/realtime/SegmentCompletionTest.java | 38 +++++++++
.../realtime/RealtimeSegmentDataManager.java | 20 ++---
.../realtime/RealtimeSegmentDataManagerTest.java | 10 ++-
7 files changed, 345 insertions(+), 23 deletions(-)
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
b/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
index 797681a64a0..2788ead958d 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/protocols/SegmentCompletionProtocol.java
@@ -26,6 +26,7 @@ import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeUnit;
+import javax.annotation.Nullable;
import org.apache.pinot.common.utils.URIUtils;
import org.apache.pinot.spi.utils.JsonUtils;
@@ -129,6 +130,7 @@ public class SegmentCompletionProtocol {
public static final String PARAM_MEMORY_USED_BYTES = "memoryUsedBytes";
public static final String PARAM_SEGMENT_SIZE_BYTES = "segmentSizeBytes";
public static final String PARAM_REASON = "reason";
+ public static final String PARAM_REASON_CODE = "reasonCode";
// Sent by servers to request additional time to build
public static final String PARAM_EXTRA_TIME_SEC = "extraTimeSec";
// Sent by servers to indicate the number of rows read so far
@@ -150,6 +152,62 @@ public class SegmentCompletionProtocol {
// (like size reaching close to its limit or number of col values for a col
is about to overflow int max)
public static final String REASON_INDEX_CAPACITY_THRESHOLD_BREACHED =
"indexCapacityThresholdBreached";
+ /// Stable wire codes for server stop reasons.
+ ///
+ /// Numeric IDs are part of the segment-completion protocol and must remain
unique, stable, and non-reusable.
+ /// Controllers ignore unknown numeric `reasonCode` values and fall back to
the retained legacy `reason` value.
+ /// When both a known `reasonCode` and a legacy `reason` are present, the
numeric code takes precedence and
+ /// normalizes the legacy reason string carried through the request params.
+ public enum ReasonCode {
+ ROW_LIMIT(100, REASON_ROW_LIMIT),
+ TIME_LIMIT(110, REASON_TIME_LIMIT),
+ END_OF_PARTITION_GROUP(120, REASON_END_OF_PARTITION_GROUP),
+ FORCE_COMMIT_MESSAGE_RECEIVED(130, REASON_FORCE_COMMIT_MESSAGE_RECEIVED),
+ INDEX_CAPACITY_THRESHOLD_BREACHED(140,
REASON_INDEX_CAPACITY_THRESHOLD_BREACHED);
+
+ private static final Map<Integer, ReasonCode> BY_ID = new HashMap<>();
+
+ static {
+ for (ReasonCode value : values()) {
+ BY_ID.put(value.getId(), value);
+ }
+ }
+
+ private final int _id;
+ private final String _reason;
+
+ ReasonCode(int id, String reason) {
+ _id = id;
+ _reason = reason;
+ }
+
+ public int getId() {
+ return _id;
+ }
+
+ public String getReason() {
+ return _reason;
+ }
+
+ @Nullable
+ public static ReasonCode fromCode(int reasonCode) {
+ return BY_ID.get(reasonCode);
+ }
+
+ @Nullable
+ public static ReasonCode fromReason(@Nullable String reason) {
+ if (reason == null) {
+ return null;
+ }
+ for (ReasonCode value : values()) {
+ if (value.getReason().equals(reason)) {
+ return value;
+ }
+ }
+ return null;
+ }
+ }
+
// Canned responses
public static final Response RESP_NOT_LEADER =
new Response(new
Response.Params().withStatus(ControllerResponseStatus.NOT_LEADER));
@@ -201,6 +259,9 @@ public class SegmentCompletionProtocol {
if (_params.getReason() != null) {
params.put(PARAM_REASON, _params.getReason());
}
+ if (_params.getReasonCode() != null) {
+ params.put(PARAM_REASON_CODE,
String.valueOf(_params.getReasonCode().getId()));
+ }
if (_params.getBuildTimeMillis() > 0) {
params.put(PARAM_BUILD_TIME_MILLIS,
String.valueOf(_params.getBuildTimeMillis()));
}
@@ -231,7 +292,10 @@ public class SegmentCompletionProtocol {
public static class Params {
private String _segmentName;
private String _instanceId;
+ @Nullable
private String _reason;
+ @Nullable
+ private ReasonCode _reasonCode;
private int _numRows;
private long _buildTimeMillis;
private long _waitTimeMillis;
@@ -253,6 +317,7 @@ public class SegmentCompletionProtocol {
_segmentSizeBytes = SEGMENT_SIZE_BYTES_DEFAULT;
_streamPartitionMsgOffset = null;
_reason = null;
+ _reasonCode = null;
}
public Params(Params params) {
@@ -267,6 +332,7 @@ public class SegmentCompletionProtocol {
_segmentSizeBytes = params.getSegmentSizeBytes();
_streamPartitionMsgOffset = params.getStreamPartitionMsgOffset();
_reason = params.getReason();
+ _reasonCode = params.getReasonCode();
}
public Params withSegmentName(String segmentName) {
@@ -279,8 +345,29 @@ public class SegmentCompletionProtocol {
return this;
}
- public Params withReason(String reason) {
+ public Params withReason(@Nullable String reason) {
_reason = reason;
+ _reasonCode = ReasonCode.fromReason(reason);
+ return this;
+ }
+
+ /// Resolves the `reasonCode` query parameter. Missing or unknown values
are ignored so the retained legacy
+ /// `reason` can continue to be used as fallback.
+ public Params withReasonCodeParam(@Nullable Integer reasonCode) {
+ ReasonCode parsedReasonCode = reasonCode != null ?
ReasonCode.fromCode(reasonCode) : null;
+ if (parsedReasonCode != null) {
+ return withReasonCode(parsedReasonCode);
+ }
+ return this;
+ }
+
+ /// Sets the typed stop reason. A known code also sets the legacy reason
string so older controllers can ignore
+ /// `reasonCode` and still use `reason`. Passing `null` clears only the
typed code.
+ public Params withReasonCode(@Nullable ReasonCode reasonCode) {
+ _reasonCode = reasonCode;
+ if (reasonCode != null) {
+ _reason = reasonCode.getReason();
+ }
return this;
}
@@ -328,10 +415,16 @@ public class SegmentCompletionProtocol {
return _segmentName;
}
+ @Nullable
public String getReason() {
return _reason;
}
+ @Nullable
+ public ReasonCode getReasonCode() {
+ return _reasonCode;
+ }
+
public String getInstanceId() {
return _instanceId;
}
@@ -372,6 +465,7 @@ public class SegmentCompletionProtocol {
return "Segment name: " + _segmentName
+ ",Instance Id: " + _instanceId
+ ",Reason: " + _reason
+ + ",ReasonCode: " + _reasonCode
+ ",NumRows: " + _numRows
+ ",BuildTimeMillis: " + _buildTimeMillis
+ ",WaitTimeMillis: " + _waitTimeMillis
diff --git
a/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
b/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
index f5854cdb060..15483b220c8 100644
---
a/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
+++
b/pinot-common/src/test/java/org/apache/pinot/common/protocols/SegmentCompletionProtocolTest.java
@@ -20,7 +20,9 @@ package org.apache.pinot.common.protocols;
import java.net.URI;
import java.util.Arrays;
+import java.util.HashSet;
import java.util.Map;
+import java.util.Set;
import java.util.stream.Collectors;
import org.apache.pinot.common.utils.URIUtils;
import org.testng.Assert;
@@ -47,6 +49,7 @@ public class SegmentCompletionProtocolTest {
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
"UNKNOWN_SEGMENT");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID),
"UNKNOWN_INSTANCE");
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON));
+
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC));
@@ -69,6 +72,7 @@ public class SegmentCompletionProtocolTest {
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
"foo__0__0__12345Z");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID),
"Server_localhost_8099");
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON));
+
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS));
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC));
@@ -79,7 +83,8 @@ public class SegmentCompletionProtocolTest {
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET));
params = new
SegmentCompletionProtocol.Request.Params().withSegmentName("foo__0__0__12345Z")
-
.withInstanceId("Server_localhost_8099").withReason("ROW_LIMIT").withBuildTimeMillis(1000)
+
.withInstanceId("Server_localhost_8099").withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+ .withBuildTimeMillis(1000)
.withWaitTimeMillis(2000).withExtraTimeSec(3000).withMemoryUsedBytes(4000).withSegmentSizeBytes(5000)
.withNumRows(6000).withSegmentLocation("/tmp/segment").withStreamPartitionMsgOffset("7000");
SegmentCompletionProtocol.SegmentCommitStartRequest
commitStartRequestWithAllParams =
@@ -89,7 +94,9 @@ public class SegmentCompletionProtocolTest {
Arrays.stream(uri.getQuery().split("&")).collect(Collectors.toMap(e ->
e.split("=")[0], e -> e.split("=")[1]));
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_SEGMENT_NAME),
"foo__0__0__12345Z");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_INSTANCE_ID),
"Server_localhost_8099");
- Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
"ROW_LIMIT");
+ Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
+ SegmentCompletionProtocol.REASON_ROW_LIMIT);
+
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE),
"100");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS),
"1000");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS),
"2000");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC),
"3000");
@@ -122,6 +129,7 @@ public class SegmentCompletionProtocolTest {
URIUtils.encode("Server_localhost_8099"));
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
"%7B%22type%22%3A%22ROW_LIMIT%22%2C%20%22value%22%3A1000%7D");
+
Assert.assertNull(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE));
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS),
"1000");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS),
"2000");
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_EXTRA_TIME_SEC),
"3000");
@@ -133,4 +141,81 @@ public class SegmentCompletionProtocolTest {
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET),
URIUtils.encode("{\"shardId-000000000001\":\"49615238429973311938200772279310862572716999467690098706\"}"));
}
+
+ @Test
+ public void testReasonCodeCompatibilityAndPrecedence() {
+ SegmentCompletionProtocol.Request.Params params =
+ new
SegmentCompletionProtocol.Request.Params().withReasonCode(SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+ Assert.assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+ Assert.assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+
+ params = new
SegmentCompletionProtocol.Request.Params().withReason("customReason")
+ .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+ Assert.assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+ Assert.assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+
+ params = new SegmentCompletionProtocol.Request.Params()
+ .withReasonCodeParam(110);
+ Assert.assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+ Assert.assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+
+ params = new
SegmentCompletionProtocol.Request.Params().withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+ .withReasonCodeParam(999);
+ Assert.assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+ Assert.assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+
+ params = new
SegmentCompletionProtocol.Request.Params().withReasonCodeParam(999);
+ Assert.assertNull(params.getReason());
+ Assert.assertNull(params.getReasonCode());
+ }
+
+ @Test
+ public void testReasonCodeValuesAreStableAndUnique() {
+ Object[][] expectedValues = {
+ {SegmentCompletionProtocol.ReasonCode.ROW_LIMIT, 100,
SegmentCompletionProtocol.REASON_ROW_LIMIT},
+ {SegmentCompletionProtocol.ReasonCode.TIME_LIMIT, 110,
SegmentCompletionProtocol.REASON_TIME_LIMIT},
+ {SegmentCompletionProtocol.ReasonCode.END_OF_PARTITION_GROUP, 120,
+ SegmentCompletionProtocol.REASON_END_OF_PARTITION_GROUP},
+ {SegmentCompletionProtocol.ReasonCode.FORCE_COMMIT_MESSAGE_RECEIVED,
130,
+ SegmentCompletionProtocol.REASON_FORCE_COMMIT_MESSAGE_RECEIVED},
+
{SegmentCompletionProtocol.ReasonCode.INDEX_CAPACITY_THRESHOLD_BREACHED, 140,
+ SegmentCompletionProtocol.REASON_INDEX_CAPACITY_THRESHOLD_BREACHED}
+ };
+
+ Set<Integer> ids = new HashSet<>();
+ for (Object[] expectedValue : expectedValues) {
+ SegmentCompletionProtocol.ReasonCode reasonCode =
+ (SegmentCompletionProtocol.ReasonCode) expectedValue[0];
+ int id = (int) expectedValue[1];
+ String reason = (String) expectedValue[2];
+
+ Assert.assertTrue(ids.add(id), "Reason code IDs must be unique");
+ Assert.assertEquals(reasonCode.getId(), id);
+ Assert.assertEquals(reasonCode.getReason(), reason);
+ Assert.assertEquals(SegmentCompletionProtocol.ReasonCode.fromCode(id),
reasonCode);
+
Assert.assertEquals(SegmentCompletionProtocol.ReasonCode.fromReason(reason),
reasonCode);
+ }
+
+ Assert.assertEquals(ids.size(),
SegmentCompletionProtocol.ReasonCode.values().length);
+ }
+
+ @Test
+ public void testReasonCodeRequestRetainsLegacyReasonForOldController()
+ throws Exception {
+ SegmentCompletionProtocol.Request.Params params = new
SegmentCompletionProtocol.Request.Params()
+ .withSegmentName("foo__0__0__12345Z")
+ .withInstanceId("Server_localhost_8099")
+ .withStreamPartitionMsgOffset("7000")
+ .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+ SegmentCompletionProtocol.SegmentConsumedRequest request =
+ new SegmentCompletionProtocol.SegmentConsumedRequest(params);
+
+ URI uri = new URI(request.getUrl("localhost:8080", "http"));
+ Map<String, String> paramsMap =
+ Arrays.stream(uri.getQuery().split("&")).collect(Collectors.toMap(e ->
e.split("=")[0], e -> e.split("=")[1]));
+
+ Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON),
+ SegmentCompletionProtocol.REASON_ROW_LIMIT);
+
Assert.assertEquals(paramsMap.get(SegmentCompletionProtocol.PARAM_REASON_CODE),
"100");
+ }
}
diff --git
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
index 4ebd5fa7de6..7f0cc8f7aec 100644
---
a/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
+++
b/pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlers.java
@@ -122,6 +122,7 @@ public class LLCSegmentCompletionHandlers {
@QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_NAME) String
segmentName,
@QueryParam(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET)
String streamPartitionMsgOffset,
@QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer
stopReasonCode,
@QueryParam(SegmentCompletionProtocol.PARAM_MEMORY_USED_BYTES) long
memoryUsedBytes,
@QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows) {
if (instanceId == null || segmentName == null || streamPartitionMsgOffset
== null) {
@@ -135,6 +136,7 @@ public class LLCSegmentCompletionHandlers {
.withSegmentName(segmentName)
.withStreamPartitionMsgOffset(streamPartitionMsgOffset)
.withReason(stopReason)
+ .withReasonCodeParam(stopReasonCode)
.withMemoryUsedBytes(memoryUsedBytes)
.withNumRows(numRows);
LOGGER.info("Processing segmentConsumed: {}", requestParams);
@@ -152,7 +154,8 @@ public class LLCSegmentCompletionHandlers {
public String
segmentStoppedConsuming(@QueryParam(SegmentCompletionProtocol.PARAM_INSTANCE_ID)
String instanceId,
@QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_NAME) String
segmentName,
@QueryParam(SegmentCompletionProtocol.PARAM_STREAM_PARTITION_MSG_OFFSET)
String streamPartitionMsgOffset,
- @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason) {
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer
stopReasonCode) {
if (instanceId == null || segmentName == null || streamPartitionMsgOffset
== null) {
LOGGER.error("Invalid call: segmentName={}, instanceId={},
streamPartitionMsgOffset={}", segmentName, instanceId,
streamPartitionMsgOffset);
@@ -163,7 +166,8 @@ public class LLCSegmentCompletionHandlers {
.withInstanceId(instanceId)
.withSegmentName(segmentName)
.withStreamPartitionMsgOffset(streamPartitionMsgOffset)
- .withReason(stopReason);
+ .withReason(stopReason)
+ .withReasonCodeParam(stopReasonCode);
LOGGER.info("Processing segmentStoppedConsuming: {}", requestParams);
String response =
_segmentCompletionManager.segmentStoppedConsuming(requestParams).toJsonString();
@@ -183,7 +187,9 @@ public class LLCSegmentCompletionHandlers {
@QueryParam(SegmentCompletionProtocol.PARAM_BUILD_TIME_MILLIS) long
buildTimeMillis,
@QueryParam(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS) long
waitTimeMillis,
@QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows,
- @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long
segmentSizeBytes) {
+ @QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long
segmentSizeBytes,
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer
stopReasonCode) {
if (instanceId == null || segmentName == null || streamPartitionMsgOffset
== null) {
LOGGER.error("Invalid call: segmentName={}, instanceId={},
streamPartitionMsgOffset={}", segmentName, instanceId,
streamPartitionMsgOffset);
@@ -198,7 +204,9 @@ public class LLCSegmentCompletionHandlers {
.withBuildTimeMillis(buildTimeMillis)
.withWaitTimeMillis(waitTimeMillis)
.withNumRows(numRows)
- .withSegmentSizeBytes(segmentSizeBytes);
+ .withSegmentSizeBytes(segmentSizeBytes)
+ .withReason(stopReason)
+ .withReasonCodeParam(stopReasonCode);
LOGGER.info("Processing segmentCommitStart: {}", requestParams);
String response =
_segmentCompletionManager.segmentCommitStart(requestParams).toJsonString();
@@ -273,7 +281,9 @@ public class LLCSegmentCompletionHandlers {
@QueryParam(SegmentCompletionProtocol.PARAM_WAIT_TIME_MILLIS) long
waitTimeMillis,
@QueryParam(SegmentCompletionProtocol.PARAM_ROW_COUNT) int numRows,
@QueryParam(SegmentCompletionProtocol.PARAM_SEGMENT_SIZE_BYTES) long
segmentSizeBytes,
- @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
FormDataMultiPart metadataFiles) {
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON) String stopReason,
+ @QueryParam(SegmentCompletionProtocol.PARAM_REASON_CODE) Integer
stopReasonCode,
+ FormDataMultiPart metadataFiles) {
if (instanceId == null || segmentName == null || segmentLocation == null
|| metadataFiles == null
|| streamPartitionMsgOffset == null) {
LOGGER.error("Invalid call: segmentName={}, instanceId={},
segmentLocation={}, streamPartitionMsgOffset={}",
@@ -292,7 +302,8 @@ public class LLCSegmentCompletionHandlers {
.withWaitTimeMillis(waitTimeMillis)
.withNumRows(numRows)
.withMemoryUsedBytes(memoryUsedBytes)
- .withReason(stopReason);
+ .withReason(stopReason)
+ .withReasonCodeParam(stopReasonCode);
LOGGER.info("Processing segmentCommitEndWithMetadata: {}", requestParams);
SegmentMetadataImpl segmentMetadata;
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
new file mode 100644
index 00000000000..83bf54cdda2
--- /dev/null
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/api/resources/LLCSegmentCompletionHandlersTest.java
@@ -0,0 +1,92 @@
+/**
+ * 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.pinot.controller.api.resources;
+
+import org.apache.pinot.common.protocols.SegmentCompletionProtocol;
+import
org.apache.pinot.controller.helix.core.realtime.SegmentCompletionManager;
+import org.mockito.ArgumentCaptor;
+import org.testng.annotations.Test;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+import static org.testng.Assert.assertEquals;
+
+
+/// Tests reason-code binding at the segment-completion REST resource boundary.
+public class LLCSegmentCompletionHandlersTest {
+ private static final String INSTANCE_ID = "Server_localhost_8099";
+ private static final String SEGMENT_NAME = "foo__0__0__12345Z";
+ private static final String OFFSET = "7000";
+
+ @Test
+ public void testSegmentCommitStartBindsReasonAndReasonCode() {
+ SegmentCompletionManager segmentCompletionManager =
mock(SegmentCompletionManager.class);
+
when(segmentCompletionManager.segmentCommitStart(any())).thenReturn(SegmentCompletionProtocol.RESP_COMMIT_CONTINUE);
+ LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+ handler._segmentCompletionManager = segmentCompletionManager;
+
+ handler.segmentCommitStart(INSTANCE_ID, SEGMENT_NAME, OFFSET, 4000, 1000,
2000, 6000, 5000, "legacyReason", 100);
+
+ SegmentCompletionProtocol.Request.Params params =
captureSegmentCommitStartParams(segmentCompletionManager);
+ assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+ assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+ }
+
+ @Test
+ public void testSegmentCommitStartFallsBackToLegacyReasonForUnknownCode() {
+ SegmentCompletionManager segmentCompletionManager =
mock(SegmentCompletionManager.class);
+
when(segmentCompletionManager.segmentCommitStart(any())).thenReturn(SegmentCompletionProtocol.RESP_COMMIT_CONTINUE);
+ LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+ handler._segmentCompletionManager = segmentCompletionManager;
+
+ handler.segmentCommitStart(INSTANCE_ID, SEGMENT_NAME, OFFSET, 4000, 1000,
2000, 6000, 5000,
+ SegmentCompletionProtocol.REASON_TIME_LIMIT, 999);
+
+ SegmentCompletionProtocol.Request.Params params =
captureSegmentCommitStartParams(segmentCompletionManager);
+ assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_TIME_LIMIT);
+ assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.TIME_LIMIT);
+ }
+
+ @Test
+ public void testSegmentConsumedPrefersKnownReasonCodeOverLegacyReason() {
+ SegmentCompletionManager segmentCompletionManager =
mock(SegmentCompletionManager.class);
+
when(segmentCompletionManager.segmentConsumed(any())).thenReturn(SegmentCompletionProtocol.RESP_FAILED);
+ LLCSegmentCompletionHandlers handler = new LLCSegmentCompletionHandlers();
+ handler._segmentCompletionManager = segmentCompletionManager;
+
+ handler.segmentConsumed(INSTANCE_ID, SEGMENT_NAME, OFFSET, "legacyReason",
100, 4000, 6000);
+
+ ArgumentCaptor<SegmentCompletionProtocol.Request.Params> paramsCaptor =
+
ArgumentCaptor.forClass(SegmentCompletionProtocol.Request.Params.class);
+ verify(segmentCompletionManager).segmentConsumed(paramsCaptor.capture());
+ SegmentCompletionProtocol.Request.Params params = paramsCaptor.getValue();
+ assertEquals(params.getReason(),
SegmentCompletionProtocol.REASON_ROW_LIMIT);
+ assertEquals(params.getReasonCode(),
SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+ }
+
+ private static SegmentCompletionProtocol.Request.Params
captureSegmentCommitStartParams(
+ SegmentCompletionManager segmentCompletionManager) {
+ ArgumentCaptor<SegmentCompletionProtocol.Request.Params> paramsCaptor =
+
ArgumentCaptor.forClass(SegmentCompletionProtocol.Request.Params.class);
+
verify(segmentCompletionManager).segmentCommitStart(paramsCaptor.capture());
+ return paramsCaptor.getValue();
+ }
+}
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
index 4c33fc249dd..95f80ee9a2a 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/realtime/SegmentCompletionTest.java
@@ -589,6 +589,44 @@ public class SegmentCompletionTest {
Assert.assertEquals(response.getStatus(),
SegmentCompletionProtocol.ControllerResponseStatus.HOLD);
}
+ @Test
+ public void testWinnerOnRowLimitReasonCodeOverridesLegacyReason() {
+ SegmentCompletionProtocol.Response response;
+ Request.Params params;
+ _segmentCompletionMgr._seconds = 10L;
+ params = new
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+ .withSegmentName(_segmentNameStr)
+ .withReason("unknownReason")
+ .withReasonCode(SegmentCompletionProtocol.ReasonCode.ROW_LIMIT);
+ response = _segmentCompletionMgr.segmentConsumed(params);
+ Assert.assertEquals(response.getStatus(),
SegmentCompletionProtocol.ControllerResponseStatus.COMMIT);
+ }
+
+ @Test
+ public void testUnknownReasonCodeFallsBackToLegacyReason() {
+ SegmentCompletionProtocol.Response response;
+ Request.Params params;
+ _segmentCompletionMgr._seconds = 10L;
+ params = new
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+ .withSegmentName(_segmentNameStr)
+ .withReason(SegmentCompletionProtocol.REASON_ROW_LIMIT)
+ .withReasonCodeParam(999);
+ response = _segmentCompletionMgr.segmentConsumed(params);
+ Assert.assertEquals(response.getStatus(),
SegmentCompletionProtocol.ControllerResponseStatus.COMMIT);
+ }
+
+ @Test
+ public void testUnknownReasonCodeWithoutLegacyReason() {
+ SegmentCompletionProtocol.Response response;
+ Request.Params params;
+ _segmentCompletionMgr._seconds = 10L;
+ params = new
Request.Params().withInstanceId(S_1).withStreamPartitionMsgOffset(_s1Offset.toString())
+ .withSegmentName(_segmentNameStr)
+ .withReasonCodeParam(999);
+ response = _segmentCompletionMgr.segmentConsumed(params);
+ Assert.assertEquals(response.getStatus(),
SegmentCompletionProtocol.ControllerResponseStatus.HOLD);
+ }
+
@Test
public void testWinnerOnForceCommit() {
SegmentCompletionProtocol.Response response;
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
index 1e239ce14da..f096ab6cc92 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
@@ -335,7 +335,7 @@ public class RealtimeSegmentDataManager extends
SegmentDataManager {
private long _consumeStartTime = -1;
private long _lastLogTime = 0;
private int _lastConsumedCount = 0;
- private String _stopReason = null;
+ private SegmentCompletionProtocol.ReasonCode _stopReasonCode = null;
private final Semaphore _segBuildSemaphore;
private final boolean _isOffHeap;
/// Whether null handling is enabled by default. This value is only used if
@@ -376,35 +376,35 @@ public class RealtimeSegmentDataManager extends
SegmentDataManager {
_startTimeMs, now, _numRowsConsumed, _numRowsIndexed);
_stopReasonPrinted = true;
}
- _stopReason = SegmentCompletionProtocol.REASON_TIME_LIMIT;
+ _stopReasonCode = SegmentCompletionProtocol.ReasonCode.TIME_LIMIT;
return true;
} else if (_numRowsIndexed >= _segmentMaxRowCount) {
_segmentLogger.info("Stopping consumption due to row limit nRows={}
numRowsIndexed={}, numRowsConsumed={}",
_segmentMaxRowCount, _numRowsIndexed, _numRowsConsumed);
- _stopReason = SegmentCompletionProtocol.REASON_ROW_LIMIT;
+ _stopReasonCode = SegmentCompletionProtocol.ReasonCode.ROW_LIMIT;
return true;
} else if (_endOfPartitionGroup) {
_segmentLogger.info("Stopping consumption due to end of
partitionGroup reached nRows={} numRowsIndexed={}, "
+ "numRowsConsumed={}", _segmentMaxRowCount, _numRowsIndexed,
_numRowsConsumed);
- _stopReason =
SegmentCompletionProtocol.REASON_END_OF_PARTITION_GROUP;
+ _stopReasonCode =
SegmentCompletionProtocol.ReasonCode.END_OF_PARTITION_GROUP;
return true;
} else if (_forceCommitMessageReceived) {
_segmentLogger.info("Stopping consumption due to force commit -
numRowsConsumed={} numRowsIndexed={}",
_numRowsConsumed, _numRowsIndexed);
- _stopReason =
SegmentCompletionProtocol.REASON_FORCE_COMMIT_MESSAGE_RECEIVED;
+ _stopReasonCode =
SegmentCompletionProtocol.ReasonCode.FORCE_COMMIT_MESSAGE_RECEIVED;
return true;
} else if (!canAddMore()) {
_segmentLogger.info(
"Stopping consumption as mutable index cannot consume more rows
- numRowsConsumed={} "
+ "numRowsIndexed={}",
_numRowsConsumed, _numRowsIndexed);
- _stopReason =
SegmentCompletionProtocol.REASON_INDEX_CAPACITY_THRESHOLD_BREACHED;
+ _stopReasonCode =
SegmentCompletionProtocol.ReasonCode.INDEX_CAPACITY_THRESHOLD_BREACHED;
return true;
}
return false;
case CATCHING_UP:
- _stopReason = null;
+ _stopReasonCode = null;
// We have posted segmentConsumed() at least once, and the controller
is asking us to catch up to a certain
// offset.
// There is no time limit here, so just check to see that we are still
within the offset we need to reach.
@@ -1071,7 +1071,7 @@ public class RealtimeSegmentDataManager extends
SegmentDataManager {
private boolean startSegmentCommit() {
SegmentCompletionProtocol.Request.Params params = new
SegmentCompletionProtocol.Request.Params();
params.withSegmentName(_segmentNameStr).withStreamPartitionMsgOffset(_currentOffset.toString())
-
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReason(_stopReason);
+
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReasonCode(_stopReasonCode);
if (_isOffHeap) {
params.withMemoryUsedBytes(_memoryManager.getTotalAllocatedBytes());
}
@@ -1373,7 +1373,7 @@ public class RealtimeSegmentDataManager extends
SegmentDataManager {
SegmentCompletionProtocol.Request.Params params = new
SegmentCompletionProtocol.Request.Params();
params.withSegmentName(_segmentNameStr).withStreamPartitionMsgOffset(_currentOffset.toString())
-
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReason(_stopReason)
+
.withNumRows(_numRowsIndexed).withInstanceId(_instanceId).withReasonCode(_stopReasonCode)
.withBuildTimeMillis(_segmentBuildDescriptor.getBuildTimeMillis())
.withSegmentSizeBytes(_segmentBuildDescriptor.getSegmentSizeBytes())
.withWaitTimeMillis(_segmentBuildDescriptor.getWaitTimeMillis());
@@ -1610,7 +1610,7 @@ public class RealtimeSegmentDataManager extends
SegmentDataManager {
// Retry maybe once if leader is not found.
SegmentCompletionProtocol.Request.Params params = new
SegmentCompletionProtocol.Request.Params();
params.withStreamPartitionMsgOffset(_currentOffset.toString()).withSegmentName(_segmentNameStr)
-
.withReason(_stopReason).withNumRows(_numRowsIndexed).withInstanceId(_instanceId);
+
.withReasonCode(_stopReasonCode).withNumRows(_numRowsIndexed).withInstanceId(_instanceId);
if (_isOffHeap) {
params.withMemoryUsedBytes(_memoryManager.getTotalAllocatedBytes());
}
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
index 7aec91dc0d2..c778794ffb8 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
@@ -1491,7 +1491,7 @@ public class RealtimeSegmentDataManagerTest {
public Field _state;
public Field _shouldStop;
- public Field _stopReason;
+ public Field _stopReasonCode;
public Field _segmentBuildFailedWithDeterministicError;
public boolean _failSegmentBuildAndReplace = false;
// When set, buildSegmentAndReplace runs the real implementation
(including the upsert CRC guard) instead of being
@@ -1545,8 +1545,8 @@ public class RealtimeSegmentDataManagerTest {
_state.setAccessible(true);
_shouldStop =
RealtimeSegmentDataManager.class.getDeclaredField("_shouldStop");
_shouldStop.setAccessible(true);
- _stopReason =
RealtimeSegmentDataManager.class.getDeclaredField("_stopReason");
- _stopReason.setAccessible(true);
+ _stopReasonCode =
RealtimeSegmentDataManager.class.getDeclaredField("_stopReasonCode");
+ _stopReasonCode.setAccessible(true);
_segmentBuildFailedWithDeterministicError =
RealtimeSegmentDataManager.class.getDeclaredField("_segmentBuildFailedWithDeterministicError");
_segmentBuildFailedWithDeterministicError.setAccessible(true);
@@ -1578,7 +1578,9 @@ public class RealtimeSegmentDataManagerTest {
public String getStopReason() {
try {
- return (String) _stopReason.get(this);
+ SegmentCompletionProtocol.ReasonCode stopReasonCode =
+ (SegmentCompletionProtocol.ReasonCode) _stopReasonCode.get(this);
+ return stopReasonCode != null ? stopReasonCode.getReason() : null;
} catch (Exception e) {
Assert.fail();
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]