Skip to content
Draft
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 @@ -82,6 +82,8 @@ public enum ControllerMeter implements AbstractMetrics.Meter {
AUDIT_REQUEST_PAYLOAD_TRUNCATED("count", true),
// Upsert compact merge task metrics
UPSERT_COMPACT_MERGE_SEGMENT_SKIPPED_CONSENSUS_FAILURE("UpsertCompactMergeSegmentsSkipped", false),
// Segments UpsertCompactionTask refused to schedule because the replicas reported different valid doc counts
UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE("UpsertCompactionSegmentsSkipped", false),
// Query workload propagation metrics
QUERY_WORKLOAD_PROPAGATION_COUNT("count", true),
QUERY_WORKLOAD_PROPAGATION_ERROR("count", true),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,8 @@
import org.apache.pinot.common.auth.AuthProviderUtils;
import org.apache.pinot.common.auth.NullAuthProvider;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.metrics.ControllerMeter;
import org.apache.pinot.common.metrics.ControllerMetrics;
import org.apache.pinot.common.restlet.resources.ValidDocIdsBitmapResponse;
import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo;
import org.apache.pinot.common.restlet.resources.ValidDocIdsType;
Expand Down Expand Up @@ -470,6 +472,18 @@ public static RoaringBitmap getValidDocIdFromServerMatchingCrc(String tableNameW
public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(String taskType,
SegmentZKMetadata segmentZKMetadata, @Nullable List<ValidDocIdsMetadataInfo> replicas, int expectedReplicaCount,
MinionConstants.ValidDocIdsConsensusMode consensusMode) {
return selectValidDocIdsMetadataForConsensus(taskType, segmentZKMetadata, replicas, expectedReplicaCount,
consensusMode, null, null);
}

/// Same, and additionally reports an EQUAL-mode disagreement on `controllerMetrics` as
/// [ControllerMeter#UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE]. Only that skip reason is metered: the
/// others mean "ask again later", not that the replicas hold different data.
@Nullable
public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(String taskType,
SegmentZKMetadata segmentZKMetadata, @Nullable List<ValidDocIdsMetadataInfo> replicas, int expectedReplicaCount,
MinionConstants.ValidDocIdsConsensusMode consensusMode, @Nullable ControllerMetrics controllerMetrics,
@Nullable String tableNameWithType) {
String segmentName = segmentZKMetadata.getSegmentName();
if (CollectionUtils.isEmpty(replicas)) {
return null;
Expand Down Expand Up @@ -534,9 +548,15 @@ public static ValidDocIdsMetadataInfo selectValidDocIdsMetadataForConsensus(Stri
// keeps the generator cheap - it avoids serializing a bitmap per replica back to the controller.
ValidDocIdsMetadataInfo first = usableReplicas.get(0);
for (int i = 1; i < usableReplicas.size(); i++) {
if (usableReplicas.get(i).getTotalValidDocs() != first.getTotalValidDocs()) {
LOGGER.warn("Replicas disagree on valid doc count for segment: {}, skipping segment for {}", segmentName,
taskType);
ValidDocIdsMetadataInfo other = usableReplicas.get(i);
if (other.getTotalValidDocs() != first.getTotalValidDocs()) {
LOGGER.warn("Replicas disagree on valid doc count for segment: {}, server {} reports {} valid docs and "
+ "server {} reports {}, skipping segment for {}", segmentName, first.getInstanceId(),
first.getTotalValidDocs(), other.getInstanceId(), other.getTotalValidDocs(), taskType);
if (controllerMetrics != null && tableNameWithType != null) {
controllerMetrics.addMeteredTableValue(tableNameWithType,
ControllerMeter.UPSERT_COMPACTION_SEGMENT_SKIPPED_CONSENSUS_FAILURE, 1L);
}
return null;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,13 @@
import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
import javax.annotation.Nullable;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.helix.task.TaskState;
import org.apache.pinot.common.exception.InvalidConfigException;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.metrics.ControllerMetrics;
import org.apache.pinot.common.restlet.resources.ValidDocIdsMetadataInfo;
import org.apache.pinot.common.restlet.resources.ValidDocIdsType;
import org.apache.pinot.controller.helix.core.PinotHelixResourceManager;
Expand Down Expand Up @@ -171,8 +173,9 @@ public List<PinotTaskConfig> generateTasks(List<TableConfig> tableConfigs) {
completedSegments.stream().collect(Collectors.toMap(SegmentZKMetadata::getSegmentName, Function.identity()));

SegmentSelectionResult segmentSelectionResult =
processValidDocIdsMetadata(taskConfigs, completedSegmentsMap, validDocIdsMetadataList,
validDocIdsMetadataResult.getSegmentToExpectedReplicaCount(), consensusMode);
processValidDocIdsMetadata(tableNameWithType, taskConfigs, completedSegmentsMap, validDocIdsMetadataList,
validDocIdsMetadataResult.getSegmentToExpectedReplicaCount(), consensusMode,
_clusterInfoAccessor.getControllerMetrics());
int skippedSegmentsCount = validDocIdsMetadataList.size()
- segmentSelectionResult.getSegmentsForCompaction().size()
- segmentSelectionResult.getSegmentsForDeletion().size();
Expand Down Expand Up @@ -217,10 +220,11 @@ public List<PinotTaskConfig> generateTasks(List<TableConfig> tableConfigs) {
}

@VisibleForTesting
public static SegmentSelectionResult processValidDocIdsMetadata(Map<String, String> taskConfigs,
Map<String, SegmentZKMetadata> completedSegmentsMap,
public static SegmentSelectionResult processValidDocIdsMetadata(String tableNameWithType,
Map<String, String> taskConfigs, Map<String, SegmentZKMetadata> completedSegmentsMap,
Map<String, List<ValidDocIdsMetadataInfo>> validDocIdsMetadataInfoMap,
Map<String, Integer> segmentToReplicaCount, MinionConstants.ValidDocIdsConsensusMode consensusMode) {
Map<String, Integer> segmentToReplicaCount, MinionConstants.ValidDocIdsConsensusMode consensusMode,
@Nullable ControllerMetrics controllerMetrics) {
double invalidRecordsThresholdPercent = Double.parseDouble(
taskConfigs.getOrDefault(UpsertCompactionTask.INVALID_RECORDS_THRESHOLD_PERCENT,
String.valueOf(DEFAULT_INVALID_RECORDS_THRESHOLD_PERCENT)));
Expand All @@ -242,7 +246,8 @@ public static SegmentSelectionResult processValidDocIdsMetadata(Map<String, Stri
List<ValidDocIdsMetadataInfo> replicas = validDocIdsMetadataInfoMap.get(segmentName);
ValidDocIdsMetadataInfo validDocIdsMetadata = MinionTaskUtils.selectValidDocIdsMetadataForConsensus(
MinionConstants.UpsertCompactionTask.TASK_TYPE, segment, replicas,
segmentToReplicaCount.getOrDefault(segmentName, replicas.size()), consensusMode);
segmentToReplicaCount.getOrDefault(segmentName, replicas.size()), consensusMode, controllerMetrics,
tableNameWithType);
if (validDocIdsMetadata == null) {
continue;
}
Expand Down
Loading
Loading