Jackie-Jiang commented on code in PR #19180:
URL: https://github.com/apache/pinot/pull/19180#discussion_r4187718866
##########
pinot-controller/src/main/java/org/apache/pinot/controller/api/resources/PinotSegmentRestletResource.java:
##########
@@ -317,6 +334,125 @@ public Map<String, Object> getSegmentMetadata(
}
}
+ @PUT
+ @Path("segments/{tableNameWithType}/{segmentName}/metadata")
+ @Authorize(targetType = TargetType.TABLE, paramName = "tableNameWithType",
action = Actions.Table.UPLOAD_SEGMENT)
+ @Authenticate(AccessType.UPDATE)
+ @Consumes(MediaType.APPLICATION_JSON)
+ @Produces(MediaType.APPLICATION_JSON)
+ @ApiOperation(value = "Update the custom map in the ZK metadata for a
segment",
+ notes = "Updates only the segment ZK metadata custom map without
uploading or refreshing the segment")
+ @ApiResponses(value = {
+ @ApiResponse(code = 200, message = "Success"),
+ @ApiResponse(code = 400, message = "Invalid table name, CRC, or custom
map modifier"),
+ @ApiResponse(code = 404, message = "Table or segment not found"),
+ @ApiResponse(code = 409, message = "Segment metadata changed
concurrently"),
+ @ApiResponse(code = 412, message = "Segment CRC does not match"),
+ @ApiResponse(code = 413, message = "Custom map modifier or resulting ZK
metadata is too large")
+ })
+ public SuccessResponse updateSegmentZKMetadataCustomMap(
+ @ApiParam(value = "Table name with type", required = true, example =
"myTable_OFFLINE")
+ @PathParam("tableNameWithType") String tableNameWithType,
+ @ApiParam(value = "Name of the segment", required = true)
@PathParam("segmentName") @Encoded String segmentName,
+ @ApiParam(value = "Expected segment CRC", required = true)
@HeaderParam(HttpHeaders.IF_MATCH)
+ String expectedCrcString,
+ @ApiParam(value = "Custom map modifier", required = true) InputStream
customMapModifierInputStream,
+ @Context HttpHeaders headers) {
+ tableNameWithType = DatabaseUtils.translateTableName(tableNameWithType,
headers);
+ segmentName = decodePathSegment(segmentName);
+ if (TableNameBuilder.getTableTypeFromTableName(tableNameWithType) == null)
{
+ throw new ControllerApplicationException(LOGGER,
+ String.format("Table type not provided with table name: %s",
tableNameWithType), Status.BAD_REQUEST);
+ }
+ String customMapModifierJson =
readCustomMapModifier(customMapModifierInputStream);
+
+ long expectedCrc;
+ try {
+ expectedCrc = Long.parseLong(expectedCrcString);
+ } catch (Exception e) {
+ throw new ControllerApplicationException(LOGGER, "Missing or invalid
If-Match segment CRC", Status.BAD_REQUEST,
+ e);
+ }
+
+ SegmentZKMetadataCustomMapModifier customMapModifier;
+ try {
+ customMapModifier = new
SegmentZKMetadataCustomMapModifier(customMapModifierJson);
+ } catch (Exception e) {
+ throw new ControllerApplicationException(LOGGER, "Invalid segment ZK
metadata custom map modifier",
+ Status.BAD_REQUEST, e);
+ }
+
+ ZNRecord segmentMetadataRecord =
+
_pinotHelixResourceManager.getSegmentMetadataZnRecord(tableNameWithType,
segmentName);
+ if (segmentMetadataRecord == null) {
+ throw new ControllerApplicationException(LOGGER,
+ String.format("Failed to find segment: %s in table: %s",
segmentName, tableNameWithType), Status.NOT_FOUND);
+ }
+ SegmentZKMetadata segmentZKMetadata = new
SegmentZKMetadata(segmentMetadataRecord);
+ if (segmentZKMetadata.getCrc() != expectedCrc) {
+ throw new ControllerApplicationException(LOGGER,
+ String.format("Segment CRC does not match for segment: %s in table:
%s", segmentName, tableNameWithType),
+ Status.PRECONDITION_FAILED);
+ }
+
segmentZKMetadata.setCustomMap(customMapModifier.modifyMap(segmentZKMetadata.getCustomMap()));
+ validateSegmentZKMetadataSize(segmentZKMetadata);
+ if
(!_pinotHelixResourceManager.updateZkMetadataWithoutDataChange(tableNameWithType,
segmentZKMetadata,
+ segmentMetadataRecord.getVersion())) {
Review Comment:
[P1] Respect the active segment upload lock
This also applies after narrowing the PR to the standalone API. With
parallel push protection enabled, ZKOperator.processExistingSegment() sets
segmentUploadStartTime and increments the metadata version before copying a
replacement into deep storage. This PUT can read that locked record, pass the
old CRC check, and increment the version again. The uploader then replaces the
deep-storage file but fails both its final metadata write and cleanup because
their expected version is stale (ZKOperator.java:346–377). This leaves the
replacement file paired with old metadata, no refresh notification, and the
upload lock still set.
Please return 409 while a segment upload lock is active, retaining the
version check to guard lock acquisition after the read. Add a regression test
that interleaves this request between the upload's lock acquisition and final
metadata write.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]