diff --git a/internal/venice-common/src/main/java/com/linkedin/venice/ConfigKeys.java b/internal/venice-common/src/main/java/com/linkedin/venice/ConfigKeys.java index adcccfd3e3d..e20db3dd6d5 100644 --- a/internal/venice-common/src/main/java/com/linkedin/venice/ConfigKeys.java +++ b/internal/venice-common/src/main/java/com/linkedin/venice/ConfigKeys.java @@ -343,7 +343,6 @@ private ConfigKeys() { */ public static final String DEFAULT_READ_STRATEGY = "default.read.strategy"; public static final String DEFAULT_OFFLINE_PUSH_STRATEGY = "default.offline.push.strategy"; - public static final String CONCURRENT_PUSH_DETECTION_STRATEGY = "concurrent.push.detection.strategy"; public static final String DEFAULT_ROUTING_STRATEGY = "default.routing.strategy"; public static final String DEFAULT_REPLICA_FACTOR = "default.replica.factor"; diff --git a/internal/venice-common/src/main/java/com/linkedin/venice/meta/ConcurrentPushDetectionStrategy.java b/internal/venice-common/src/main/java/com/linkedin/venice/meta/ConcurrentPushDetectionStrategy.java deleted file mode 100644 index f7cdc4d0dbd..00000000000 --- a/internal/venice-common/src/main/java/com/linkedin/venice/meta/ConcurrentPushDetectionStrategy.java +++ /dev/null @@ -1,17 +0,0 @@ -package com.linkedin.venice.meta; - -public enum ConcurrentPushDetectionStrategy { - TOPIC_BASED_ONLY(true), // Current mode using parent VT - DUAL(true), // Read from parent zk status and Pubsub and compare - PARENT_VERSION_STATUS_ONLY(false); - - private final boolean topicWriteNeeded; - - ConcurrentPushDetectionStrategy(boolean topicWriteNeeded) { - this.topicWriteNeeded = topicWriteNeeded; - } - - public boolean isTopicWriteNeeded() { - return topicWriteNeeded; - } -} diff --git a/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/TestParentControllerWithMultiDataCenter.java b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/TestParentControllerWithMultiDataCenter.java index 623405dd32f..570f8be592f 100644 --- a/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/TestParentControllerWithMultiDataCenter.java +++ b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/TestParentControllerWithMultiDataCenter.java @@ -1301,7 +1301,7 @@ public void testCompliancePushCannotKillAnotherCompliancePush() { "Second compliance push should fail because system pushes cannot kill other system pushes"); Assert.assertTrue( compliancePushResponse2.getError().contains("An ongoing push") && compliancePushResponse2.getError() - .contains("is found and it must be terminated before another push can be started"), + .contains("is still in progress and must complete or be terminated before another push can be started"), "Error should indicate an ongoing push must be terminated: " + compliancePushResponse2.getError()); } } @@ -1334,7 +1334,7 @@ public void testCompliancePushCannotKillUserPush() { "Compliance push should fail because system pushes cannot kill user pushes"); Assert.assertTrue( compliancePushResponse.getError().contains("An ongoing push") && compliancePushResponse.getError() - .contains("is found and it must be terminated before another push can be started"), + .contains("is still in progress and must complete or be terminated before another push can be started"), "Error should indicate an ongoing push must be terminated: " + compliancePushResponse.getError()); } } diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java index 76cf5bd78f7..ee1cfa305c6 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java @@ -14,7 +14,6 @@ import com.linkedin.venice.exceptions.VeniceNoClusterException; import com.linkedin.venice.hooks.StoreLifecycleHooks; import com.linkedin.venice.hooks.StoreVersionLifecycleEventOutcome; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.LifecycleHooksRecord; import com.linkedin.venice.meta.ReadWriteStoreRepository; import com.linkedin.venice.meta.Store; @@ -1278,17 +1277,6 @@ public void updateStore(String clusterName, String storeName, VersionStatus stat store.updateVersionStatus(targetVersionNum, status); if (status == ONLINE || status == PARTIALLY_ONLINE) { store.setCurrentVersion(targetVersionNum); - - // For jobs that stop polling early or for pushes that don't poll (empty push), we need to truncate the parent - // VT here to unblock the next push - String kafkaTopicName = Version.composeKafkaTopic(storeName, targetVersionNum); - ConcurrentPushDetectionStrategy strategy = - veniceControllerMultiClusterConfig.getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); - // skip truncating if the topic was not created based on ConcurrentPushDetectionStrategy - if (strategy.isTopicWriteNeeded() && !veniceParentHelixAdmin.isTopicTruncated(kafkaTopicName)) { - LOGGER.info("Truncating parent VT for {}", kafkaTopicName); - veniceParentHelixAdmin.truncateKafkaTopic(Version.composeKafkaTopic(storeName, targetVersionNum)); - } } repository.updateStore(store); } catch (Exception e) { diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java index 68676953c58..f2593951949 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java @@ -15,7 +15,6 @@ import com.linkedin.venice.controller.grpc.server.interceptor.ControllerGrpcSslSessionInterceptor; import com.linkedin.venice.controller.grpc.server.interceptor.ParentControllerRegionValidationInterceptor; import com.linkedin.venice.controller.kafka.TopicCleanupService; -import com.linkedin.venice.controller.kafka.TopicCleanupServiceForParentController; import com.linkedin.venice.controller.lingeringjob.HeartbeatCheckerOtelMetricEntity; import com.linkedin.venice.controller.server.AdminSparkServer; import com.linkedin.venice.controller.server.VeniceControllerGrpcServiceImpl; @@ -294,23 +293,12 @@ AdminSparkServer createAdminServer(boolean secure) { private TopicCleanupService createTopicCleanupService() { Admin admin = controllerService.getVeniceHelixAdmin(); - - if (multiClusterConfigs.isParent()) { - // TODO: Remove the following once ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY is fully rolled out - return new TopicCleanupServiceForParentController( - admin, - multiClusterConfigs, - pubSubTopicRepository, - new TopicCleanupServiceStats(metricsRepository), - pubSubClientsFactory); - } else { - return new TopicCleanupService( - admin, - multiClusterConfigs, - pubSubTopicRepository, - new TopicCleanupServiceStats(metricsRepository), - pubSubClientsFactory); - } + return new TopicCleanupService( + admin, + multiClusterConfigs, + pubSubTopicRepository, + new TopicCleanupServiceStats(metricsRepository), + pubSubClientsFactory); } private Optional createStoreBackupVersionCleanupService() { diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerClusterConfig.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerClusterConfig.java index 78f1b37b361..428f4dcd9be 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerClusterConfig.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerClusterConfig.java @@ -27,7 +27,6 @@ import static com.linkedin.venice.ConfigKeys.CLUSTER_TO_D2; import static com.linkedin.venice.ConfigKeys.CLUSTER_TO_SERVER_D2; import static com.linkedin.venice.ConfigKeys.CONCURRENT_INIT_ROUTINES_ENABLED; -import static com.linkedin.venice.ConfigKeys.CONCURRENT_PUSH_DETECTION_STRATEGY; import static com.linkedin.venice.ConfigKeys.CONTROLLER_ADMIN_GRPC_PORT; import static com.linkedin.venice.ConfigKeys.CONTROLLER_ADMIN_SECURE_GRPC_PORT; import static com.linkedin.venice.ConfigKeys.CONTROLLER_AUTO_MATERIALIZE_DAVINCI_PUSH_STATUS_SYSTEM_STORE; @@ -249,7 +248,6 @@ import com.linkedin.venice.controllerapi.ControllerRoute; import com.linkedin.venice.exceptions.ConfigurationException; import com.linkedin.venice.exceptions.VeniceException; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.OfflinePushStrategy; import com.linkedin.venice.meta.PersistenceType; import com.linkedin.venice.meta.ReadStrategy; @@ -512,7 +510,6 @@ public class VeniceControllerClusterConfig { private final PersistenceType persistenceType; private final ReadStrategy readStrategy; - private final ConcurrentPushDetectionStrategy concurrentPushDetectionStrategy; private final OfflinePushStrategy offlinePushStrategy; private final RoutingStrategy routingStrategy; private final int replicationFactor; @@ -785,13 +782,6 @@ public VeniceControllerClusterConfig(VeniceProperties props) { this.isSkipHybridStoreRTTopicCompactionPolicyUpdateEnabled = props.getBoolean(SKIP_HYBRID_STORE_RT_TOPIC_COMPACTION_POLICY_UPDATE_ENABLED, false); - if (props.containsKey(CONCURRENT_PUSH_DETECTION_STRATEGY)) { - this.concurrentPushDetectionStrategy = - ConcurrentPushDetectionStrategy.valueOf(props.getString(CONCURRENT_PUSH_DETECTION_STRATEGY)); - } else { - this.concurrentPushDetectionStrategy = ConcurrentPushDetectionStrategy.DUAL; - } - if (props.containsKey(DEFAULT_READ_STRATEGY)) { this.readStrategy = ReadStrategy.valueOf(props.getString(DEFAULT_READ_STRATEGY)); } else { @@ -1472,10 +1462,6 @@ public ReadStrategy getReadStrategy() { return readStrategy; } - public ConcurrentPushDetectionStrategy getConcurrentPushDetectionStrategy() { - return concurrentPushDetectionStrategy; - } - public OfflinePushStrategy getOfflinePushStrategy() { return offlinePushStrategy; } diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceHelixAdmin.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceHelixAdmin.java index e387cff228f..0a688be2893 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceHelixAdmin.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceHelixAdmin.java @@ -3115,7 +3115,7 @@ private Pair addVersion( version.getNumber(), store.getStoreLifecycleHooks()); long createBatchTopicStartTime = System.currentTimeMillis(); - if (clusterConfig.getConcurrentPushDetectionStrategy().isTopicWriteNeeded() || !isParent()) { + if (!isParent()) { topicToCreationTime.computeIfAbsent(version.kafkaTopicName(), topic -> System.currentTimeMillis()); createBatchTopics( version, @@ -4121,18 +4121,6 @@ public void deleteOneStoreVersion(String clusterName, String storeName, int vers deleteOneStoreVersion(clusterName, storeName, versionNumber, false); } - /** - * Check if we should skip truncating topic. If it's parent fabrics and the topic write is NOT needed, return true; - * Otherwise, return false. - * @param clusterName the cluster name to check - * @return true if topic truncation should be skipped, false otherwise - */ - public boolean shouldSkipTruncatingTopic(String clusterName) { - return isParent() && !getMultiClusterConfigs().getControllerConfig(clusterName) - .getConcurrentPushDetectionStrategy() - .isTopicWriteNeeded(); - } - private void deleteOneStoreVersion(String clusterName, String storeName, int versionNumber, boolean isForcedDelete) { HelixVeniceClusterResources resources = getHelixVeniceClusterResources(clusterName); try (AutoCloseableLock ignore = resources.getClusterLockManager().createStoreWriteLock(storeName)) { @@ -4169,8 +4157,8 @@ private void deleteOneStoreVersion(String clusterName, String storeName, int ver // Not using deletedVersion.get().kafkaTopicName() because it's incorrect for Zk shared stores. String versionTopicName = Version.composeKafkaTopic(storeName, deletedVersion.get().getNumber()); - // skip truncating topic if it's parent controller and topic write is not needed - if (!shouldSkipTruncatingTopic(clusterName)) { + // skip truncating topic if it's a parent controller (parent controllers do not write version topics) + if (!isParent()) { if (fatalDataValidationFailureRetentionMs != -1 && hasFatalDataValidationError) { truncateKafkaTopic(versionTopicName, fatalDataValidationFailureRetentionMs); } else { diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceParentHelixAdmin.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceParentHelixAdmin.java index 26934035790..b05b4dbdb1f 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceParentHelixAdmin.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceParentHelixAdmin.java @@ -117,7 +117,6 @@ import com.linkedin.venice.helix.Replica; import com.linkedin.venice.helix.StoragePersonaRepository; import com.linkedin.venice.helix.ZkStoreConfigAccessor; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.DegradedDcInfo; import com.linkedin.venice.meta.ETLStoreConfig; import com.linkedin.venice.meta.IngestionPauseMode; @@ -132,7 +131,6 @@ import com.linkedin.venice.meta.StoreDataAudit; import com.linkedin.venice.meta.StoreGraveyard; import com.linkedin.venice.meta.StoreInfo; -import com.linkedin.venice.meta.StoreVersionInfo; import com.linkedin.venice.meta.VeniceETLStrategy; import com.linkedin.venice.meta.Version; import com.linkedin.venice.meta.VersionStatus; @@ -210,7 +208,6 @@ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.Function; -import java.util.stream.Collectors; import javax.annotation.Nonnull; import org.apache.avro.Schema; import org.apache.commons.lang.StringUtils; @@ -234,7 +231,6 @@ public class VeniceParentHelixAdmin implements Admin { private static final StackTraceElement[] EMPTY_STACK_TRACE = new StackTraceElement[0]; - private static final long TOPIC_DELETION_DELAY_MS = 5 * Time.MS_PER_MINUTE; public static final List> RETRY_FAILURE_TYPES = Collections.singletonList(Exception.class); private static final int ROLL_FORWARD_REQUEST_TIMEOUT = 60 * Time.MS_PER_SECOND; private static final int CONTROLLER_STORE_POLL_TIMEOUT = 5 * Time.MS_PER_SECOND; @@ -270,15 +266,11 @@ public class VeniceParentHelixAdmin implements Admin { private ParentHelixOfflinePushAccessor offlinePushAccessor; /** - * Here is the way how Parent Controller is keeping errored topics when {@link #maxErroredTopicNumToKeep} > 0: - * 1. For errored topics, {@link #getOffLineJobStatus} won't truncate them; - * 2. For errored topics, {@link #killOfflinePush(String, String, boolean)} won't truncate them; - * 3. {@link #getTopicForCurrentPushJob(String, String, boolean, boolean)} will truncate the errored topics based on - * {@link #maxErroredTopicNumToKeep}; - * - * It means error topic retiring is only be triggered by next push. - * - * When {@link #maxErroredTopicNumToKeep} is 0, errored topics will be truncated right away when job is finished. + * The parent controller no longer writes or truncates version topics, so errored version topics are + * neither retained nor cleaned up here. {@link #maxErroredTopicNumToKeep} now only gates: cleanup of + * the corresponding stream-reprocessing topic in {@link #killOfflinePush(String, String, boolean)} + * (when it is 0); and, in {@link #truncateTopicsOptionally}, whether an errored push's stream-reprocessing + * topic is skipped for truncation and reported as such in the status details (when it is greater than 0). */ private int maxErroredTopicNumToKeep; @@ -1257,227 +1249,18 @@ void cleanupHistoricalVersions(String clusterName, String storeName) { } } - /** - * Check whether any topic for this store exists or not. - * The existing topic could be introduced by two cases: - * 1. The previous job push is still running; - * 2. The previous job push fails to delete this topic; - * - * For the 1st case, it is expected to refuse the new data push, - * and for the 2nd case, customer should reach out Venice team to fix this issue for now. - **/ - List existingVersionTopicsForStore(String storeName) { - List outputList = new ArrayList<>(); - TopicManager topicManager = getTopicManager(); - Set topics = topicManager.listTopics(); - String storeNameForCurrentTopic; - for (PubSubTopic topic: topics) { - if (AdminTopicUtils.isAdminTopic(topic.getName()) || AdminTopicUtils.isKafkaInternalTopic(topic.getName()) - || topic.isRealTime() || VeniceView.isViewTopic(topic.getName())) { - continue; - } - try { - storeNameForCurrentTopic = Version.parseStoreFromKafkaTopicName(topic.getName()); - } catch (Exception e) { - LOGGER.debug("Failed to parse StoreName from topic: {}, and error message: {}", topic, e.getMessage()); - continue; - } - if (storeNameForCurrentTopic.equals(storeName)) { - outputList.add(topic); - } - } - return outputList; - } - - /** - * Get the version topics list for the specified store in freshness order; the first - * topic in the list is the latest topic and the last topic is the oldest one. - * @param storeName - * @return the version topics in freshness order - */ - List getKafkaTopicsByAge(String storeName) { - List existingTopics = existingVersionTopicsForStore(storeName); - if (!existingTopics.isEmpty()) { - existingTopics.sort((t1, t2) -> { - int v1 = Version.parseVersionFromKafkaTopicName(t1.getName()); - int v2 = Version.parseVersionFromKafkaTopicName(t2.getName()); - return v2 - v1; - }); - } - return existingTopics; - } - /** * If there is no ongoing push for specified store currently, this function will return {@link Optional#empty()}, - * else will return the ongoing Kafka topic. It will also try to clean up legacy topics. + * else will return the ongoing Kafka topic. + * + *

Note: {@code isIncrementalPush} and {@code isRepush} are accepted for interface/call-site compatibility but + * are not used by this parent-controller implementation, which relies solely on parent version status. */ - public Optional getTopicForCurrentPushJob( String clusterName, String storeName, boolean isIncrementalPush, boolean isRepush) { - VeniceControllerClusterConfig controllerConfig = - getVeniceHelixAdmin().getHelixVeniceClusterResources(clusterName).getConfig(); - ConcurrentPushDetectionStrategy pushDetectionStrategy = controllerConfig.getConcurrentPushDetectionStrategy(); - if (ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY.equals(pushDetectionStrategy)) { - return getTopicForCurrentPushJobTopicBasedTracking(clusterName, storeName, isIncrementalPush, isRepush); - } else if (ConcurrentPushDetectionStrategy.DUAL.equals(pushDetectionStrategy)) { - Optional topicBased = - getTopicForCurrentPushJobTopicBasedTracking(clusterName, storeName, isIncrementalPush, isRepush); - Optional versionStatusBased = - getTopicForCurrentPushJobParentVersionStatusBasedTracking(clusterName, storeName); - if (!topicBased.equals(versionStatusBased)) { - LOGGER.error( - "getTopicForCurrentPushJob returns different value for store {} in cluster {}, topicBased: {}, versionStatusBased: {}", - storeName, - clusterName, - topicBased, - versionStatusBased); - } - return topicBased; - } else { - return getTopicForCurrentPushJobParentVersionStatusBasedTracking(clusterName, storeName); - } - } - - Optional getTopicForCurrentPushJobTopicBasedTracking( - String clusterName, - String storeName, - boolean isIncrementalPush, - boolean isRepush) { - // The first/last topic in the list is the latest/oldest version topic - List versionTopics = getKafkaTopicsByAge(storeName); - Optional latestTopic = Optional.empty(); - if (!versionTopics.isEmpty()) { - latestTopic = Optional.of(versionTopics.get(0)); - } - - if (latestTopic.isPresent()) { - LOGGER.debug("Latest kafka topic for store: {} is {}", storeName, latestTopic.get()); - final String latestTopicName = latestTopic.get().getName(); - int versionNumber = Version.parseVersionFromKafkaTopicName(latestTopicName); - Store store = getStore(clusterName, storeName); - Version version = store.getVersion(versionNumber); - boolean onlyDeferredSwap = version.isVersionSwapDeferred() && StringUtils.isEmpty(version.getTargetSwapRegion()); - boolean isTargetRegionPushWithDeferredSwap = - version != null && version.isVersionSwapDeferred() && StringUtils.isNotEmpty(version.getTargetSwapRegion()); - - if (onlyDeferredSwap) { - if (version.getStatus() == STARTED || version.getStatus() == PUSHED) { - LOGGER.error( - "Future version {} exists for store {}, please wait till the future version is made current.", - versionNumber, - storeName); - return Optional.of(latestTopic.get().getName()); - } else if (version.getStatus() == ONLINE) { - // for only deferred swap, users need to rollforward to mark it current for online status - boolean validatedChildVersion = validateChildCurrentVersions(clusterName, storeName, versionNumber); - if (!validatedChildVersion) { - return Optional.of(latestTopic.get().getName()); - } - } - } else if (isTargetRegionPushWithDeferredSwap) { - LOGGER.error( - "Future version {} exists for store {}, please wait till the future version is made current.", - versionNumber, - storeName); - return Optional.of(latestTopic.get().getName()); - } - - if (!isTopicTruncated(latestTopicName)) { - /** - * Check whether the corresponding version exists or not, since it is possible that last push - * meets Kafka topic creation timeout. - * When Kafka topic creation timeout happens, topic/job could be still running, but the version - * should not exist according to the logic in {@link VeniceHelixAdmin#addVersion}. - * However, it is possible that a different request enters this code section when the topic has been created but - * either the version information has not been persisted to Zk or the in-memory Store object. In this case, it - * is desirable to add a delay to topic deletion. - * - * If the corresponding version doesn't exist, this function will issue command to kill job to deprecate - * the incomplete topic/job. - */ - StoreVersionInfo storeVersionPair = - getVeniceHelixAdmin().waitVersion(clusterName, storeName, versionNumber, Duration.ofSeconds(30)); - if (storeVersionPair.getVersion() == null) { - // TODO: Guard this topic deletion code using a store-level lock instead. - Long inMemoryTopicCreationTime = getVeniceHelixAdmin().getInMemoryTopicCreationTime(latestTopicName); - if (inMemoryTopicCreationTime != null - && SystemTime.INSTANCE.getMilliseconds() < (inMemoryTopicCreationTime + TOPIC_DELETION_DELAY_MS)) { - throw new VeniceException( - "Failed to get version information but the topic exists and has been created recently. Try again after some time."); - } - - killOfflinePush(clusterName, latestTopicName, true); - LOGGER.info("Found topic: {} without the corresponding version, will kill it", latestTopicName); - return Optional.empty(); - } - - /** - * If Parent Controller could not infer the job status from topic retention policy, it will check the actual - * job status by sending requests to each individual datacenter. - * If the job is still running, Parent Controller will block current push. - */ - final long SLEEP_MS_BETWEEN_RETRY = TimeUnit.SECONDS.toMillis(10); - ExecutionStatus jobStatus = ExecutionStatus.PROGRESS; - Map extraInfo = new HashMap<>(); - - int retryTimes = 5; - int current = 0; - while (current++ < retryTimes) { - OfflinePushStatusInfo offlineJobStatus = getOffLinePushStatus(clusterName, latestTopicName); - jobStatus = offlineJobStatus.getExecutionStatus(); - extraInfo = offlineJobStatus.getExtraInfo(); - if (!extraInfo.containsValue(ExecutionStatus.UNKNOWN.toString())) { - break; - } - // Retry since there is a connection failure when querying job status against each datacenter - try { - timer.sleep(SLEEP_MS_BETWEEN_RETRY); - } catch (InterruptedException e) { - currentThread().interrupt(); - throw new VeniceException( - "Received InterruptedException during sleep between 'getOffLinePushStatus' calls"); - } - } - if (extraInfo.containsValue(ExecutionStatus.UNKNOWN.toString())) { - // TODO: Do we need to throw exception here?? - LOGGER.error( - "Failed to get job status for topic: {} after retrying {} times, extra info: {}", - latestTopicName, - retryTimes, - extraInfo); - } - if (!jobStatus.isTerminal()) { - LOGGER.info( - "Job status: {} for Kafka topic: {} is not terminal, extra info: {}", - jobStatus, - latestTopicName, - extraInfo); - if (latestTopic.isPresent()) { - return Optional.of(latestTopic.get().getName()); - } - return Optional.empty(); - } else { - /** - * If the job status of latestKafkaTopic is terminal and it is not an incremental push, - * it will be truncated in {@link #getOffLinePushStatus(String, String)}. - */ - if (!isIncrementalPush) { - Map currentVersionsMap = getCurrentVersionsForMultiColos(clusterName, storeName); - truncateTopicsBasedOnMaxErroredTopicNumToKeep( - versionTopics.stream().map(vt -> vt.getName()).collect(Collectors.toList()), - isRepush, - currentVersionsMap); - } - } - } - } - return Optional.empty(); - } - - Optional getTopicForCurrentPushJobParentVersionStatusBasedTracking(String clusterName, String storeName) { Store store = getStore(clusterName, storeName); if (store == null) { return Optional.empty(); @@ -1485,10 +1268,20 @@ Optional getTopicForCurrentPushJobParentVersionStatusBasedTracking(Strin int lastVersionNum = store.getLargestUsedVersionNumber(); Version lastVersion = store.getVersion(lastVersionNum); - if (lastVersionNum == NON_EXISTING_VERSION || lastVersion == null) { + if (lastVersionNum == NON_EXISTING_VERSION) { LOGGER.info("Store {} does not have any version", storeName); return Optional.empty(); } + if (lastVersion == null) { + // largestUsedVersionNumber is monotonic and isn't decremented on deleteVersion, and can also be + // advanced before the corresponding Version is added (e.g. initiateDataRecovery); in either case + // there is no ongoing push to wait on for this (missing) version. + LOGGER.info( + "Store {} largest used version {} does not correspond to an existing version; no ongoing push to wait on", + storeName, + lastVersionNum); + return Optional.empty(); + } // Terminal statuses for the latest version mean there is no in-flight push to wait on, so // the next push may proceed: @@ -1498,7 +1291,9 @@ Optional getTopicForCurrentPushJobParentVersionStatusBasedTracking(Strin // checkRollbackOriginVersionCapacityForNewPush. // - PARTIALLY_ONLINE: parent-only terminal state from a region-filtered rollback (some // regions rolled back, some didn't); same retention window applies via the same guard. - // Non-terminal statuses fall through to the polling branch below to wait on the in-flight push. + // Non-terminal statuses are resolved below: a version that defers its swap, or one in CREATED/PUSHED, + // blocks the next push outright; a non-deferred ONLINE version has no ongoing push; STARTED polls the + // child job status to decide. switch (lastVersion.getStatus()) { case KILLED: case ERROR: @@ -1520,26 +1315,49 @@ Optional getTopicForCurrentPushJobParentVersionStatusBasedTracking(Strin storeName, lastVersionNum); Optional latestTopic = Optional.of(Version.composeKafkaTopic(storeName, lastVersionNum)); - boolean onlyDeferredSwap = - lastVersion.isVersionSwapDeferred() && StringUtils.isEmpty(lastVersion.getTargetSwapRegion()); - if ((lastVersion.getStatus() == STARTED || lastVersion.getStatus() == PUSHED - || lastVersion.getStatus() == CREATED)) { - LOGGER.error( - "The push for version {} of store {} is not completed, please wait till the push is completed.", + if (lastVersion.isVersionSwapDeferred()) { + // A version that defers its swap occupies the store until it is rolled forward and current in + // every region; block the next push until then, whatever the push status. This is decided before + // the child-status poll below, whose side effects would otherwise advance the version to a + // terminal status and incorrectly unblock the next push -- for a deferred swap the push is + // "complete" yet the version is deliberately not made current. + if (validateChildCurrentVersions(clusterName, storeName, lastVersion)) { + return Optional.empty(); + } + return latestTopic; + } + + if (lastVersion.getStatus() == CREATED || lastVersion.getStatus() == PUSHED) { + // CREATED: the version exists but its push has not begun. PUSHED: a target-region push that has + // completed in its target region but not yet in the rest. In both cases a future version already + // occupies the store, so the next push must wait. STARTED is excluded here: a STARTED version may + // already be terminal in child regions but not yet reflected on the parent, so it is handled by + // the polling branch below. (Unlike STARTED, polling PUSHED children would incorrectly report the + // push as terminal, since the children have already completed it.) + LOGGER.info( + "The push for version {} (pushJobId {}) of store {} is not completed (status {}); the next push must wait.", lastVersionNum, - storeName); + lastVersion.getPushJobId(), + storeName, + lastVersion.getStatus()); return latestTopic; - } else if (onlyDeferredSwap && lastVersion.getStatus() == ONLINE) { - // for only deferred swap, users need to rollforward to mark it current for online status - boolean validateChildCurrentVersions = validateChildCurrentVersions(clusterName, storeName, lastVersionNum); - if (!validateChildCurrentVersions) { - return latestTopic; - } + } + + if (lastVersion.getStatus() == ONLINE) { + // A non-deferred ONLINE version has completed its push and is serving reads, so there is no + // ongoing push to wait on. Child regions may legitimately diverge for an ONLINE version -- e.g. a + // version deleted or rolled back in one region -- which is a stale-store condition, not an + // in-flight push, so it must not be re-polled (doing so would misreport it as in progress). + return Optional.empty(); } /** - * If the job is still running, Parent Controller will block current push. + * A STARTED version can reflect either a push that is still in flight or one that already completed + * in the child regions but whose parent version status has not yet been advanced -- the parent + * version only transitions out of STARTED when a job-status poll observes a terminal child status + * (see {@link #getOffLineJobStatus}). Poll the child job status to tell these apart: block while it + * is non-terminal, otherwise let the next push proceed. */ final long SLEEP_MS_BETWEEN_RETRY = TimeUnit.SECONDS.toMillis(10); ExecutionStatus jobStatus = ExecutionStatus.PROGRESS; @@ -1569,79 +1387,27 @@ Optional getTopicForCurrentPushJobParentVersionStatusBasedTracking(Strin return Optional.empty(); } - private boolean validateChildCurrentVersions(String clusterName, String storeName, int lastVersionNum) { + private boolean validateChildCurrentVersions(String clusterName, String storeName, Version lastVersion) { + int lastVersionNum = lastVersion.getNumber(); Map currentVersionsMap = getCurrentVersionsForMultiColos(clusterName, storeName); - for (Map.Entry entry: currentVersionsMap.entrySet()) { + List regionsNotYetCurrent = new ArrayList<>(); + for (Map.Entry entry: currentVersionsMap.entrySet()) { if (!entry.getValue().equals(lastVersionNum)) { - LOGGER.error( - "Future version {} exists for store {}, but in region {} current version {}, please wait till the future version is made current.", - lastVersionNum, - storeName, - entry.getKey(), - entry.getValue()); - return false; + regionsNotYetCurrent.add(entry.getKey()); } } - return true; - } - - /** - * Only keep {@link #maxErroredTopicNumToKeep} non-truncated topics ordered by version. It works as a general method - * for cleaning up leaking topics. ({@link #maxErroredTopicNumToKeep} is always 0.) - */ - void truncateTopicsBasedOnMaxErroredTopicNumToKeep( - List topics, - boolean isRepush, - Map currentVersionsMap) { - // Based on current logic, only 'errored' topics were not truncated. - List sortedNonTruncatedTopics = - topics.stream().filter(topic -> !isTopicTruncated(topic)).sorted((t1, t2) -> { - int v1 = Version.parseVersionFromKafkaTopicName(t1); - int v2 = Version.parseVersionFromKafkaTopicName(t2); - return v1 - v2; - }).collect(Collectors.toList()); - Set streamReprocessingTopics = - sortedNonTruncatedTopics.stream().filter(Version::isStreamReprocessingTopic).collect(Collectors.toSet()); - List sortedNonTruncatedVersionTopics = sortedNonTruncatedTopics.stream() - .filter(topic -> !Version.isStreamReprocessingTopic(topic)) - .collect(Collectors.toList()); - if (sortedNonTruncatedVersionTopics.size() <= maxErroredTopicNumToKeep) { + if (!regionsNotYetCurrent.isEmpty()) { LOGGER.info( - "Non-truncated version topics size: {} isn't bigger than maxErroredTopicNumToKeep: {}, so no topic " - + "will be truncated this time", - sortedNonTruncatedTopics.size(), - maxErroredTopicNumToKeep); - return; - } - int topicNumToTruncate = sortedNonTruncatedVersionTopics.size() - maxErroredTopicNumToKeep; - int truncatedTopicCnt = 0; - for (String topic: sortedNonTruncatedVersionTopics) { - /** - * If Venice repush somehow failed and we delete the version topic for the current version here, future incremental - * pushes will fail; therefore, keep Venice repush transparent and don't delete any VTs; future regular batch pushes - * from users will delete the VT we retain here. - * Potential improvement: After the Venice repush completes, we can automatically deletes VT from previous version, - * at the risk of not being able to roll back to previous version though, so not recommend to do such automation. - */ - if (isRepush && currentVersionsMap.containsValue(Version.parseVersionFromVersionTopicName(topic))) { - LOGGER.info( - "Do not delete the current version topic: {} since the incoming push is a Venice internal re-push.", - topic); - continue; - } - if (++truncatedTopicCnt > topicNumToTruncate) { - break; - } - truncateKafkaTopic(topic); - LOGGER.info("Errored topic: {} got truncated", topic); - String correspondingStreamReprocessingTopic = Version.composeStreamReprocessingTopicFromVersionTopic(topic); - if (streamReprocessingTopics.contains(correspondingStreamReprocessingTopic)) { - truncateKafkaTopic(correspondingStreamReprocessingTopic); - LOGGER.info( - "Corresponding stream reprocessing topic: {} also got truncated.", - correspondingStreamReprocessingTopic); - } + "Store {} version {} (status {}) defers its version swap and is not yet current in all regions; " + + "regions not yet current: {}; current version per region: {}", + storeName, + lastVersionNum, + lastVersion.getStatus(), + regionsNotYetCurrent, + currentVersionsMap); + return false; } + return true; } /** @@ -1860,11 +1626,23 @@ public Version incrementVersionIdempotent( pushJobId, storeName); } else { - String msg = version.isVersionSwapDeferred() - ? ". There is already a future version " + version.getNumber() + " exists for the store " + storeName - + " please make that version current before starting a next push." - : ". An ongoing push with pushJobId " + existingPushJobId + " and topic " + currentPushTopic.get() - + " is found and it must be terminated before another push can be started."; + String msg; + if (version.isVersionSwapDeferred() && version.getStatus() == PUSHED) { + // The blocking version finished its push but is mid deferred (colo-by-colo) version swap: + // PUSHED on the parent (push complete but swap not yet done) and not yet current in every region. + // Surface the per-region status so operators do not mistake this for an in-flight concurrent push. + Map currentVersions = getCurrentVersionsForMultiColos(clusterName, storeName); + String targetSwapRegion = version.getTargetSwapRegion(); + msg = ". Version " + version.getNumber() + " of store " + storeName + + " has completed its push and is waiting on deferred version swap (colo-by-colo roll-forward);" + + " current version per region: " + currentVersions + + (StringUtils.isNotEmpty(targetSwapRegion) ? ", target swap region(s): " + targetSwapRegion : "") + + ". Roll it forward to make it current in all regions (or roll it back) before starting a new push."; + } else { + msg = ". An ongoing push for version " + version.getNumber() + " (status " + version.getStatus() + + ", pushJobId " + existingPushJobId + ", topic " + currentPushTopic.get() + + ") is still in progress and must complete or be terminated before another push can be started."; + } VeniceException e = new ConcurrentBatchPushException( "Unable to start the push with pushJobId " + pushJobId + " for store " + storeName + msg); e.setStackTrace(EMPTY_STACK_TRACE); @@ -2349,17 +2127,6 @@ public void rollForwardToFutureVersion(String clusterName, String storeName, Str } String kafkaTopic = Version.composeKafkaTopic(storeName, futureVersionBeforeRollForward); - Version futureVersion = getStore(clusterName, storeName).getVersion(futureVersionBeforeRollForward); - boolean onlyDeferredSwap = - futureVersion.isVersionSwapDeferred() && StringUtils.isEmpty(futureVersion.getTargetSwapRegion()); - ConcurrentPushDetectionStrategy concurrentPushDetectionStrategy = - getMultiClusterConfigs().getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); - if (onlyDeferredSwap && concurrentPushDetectionStrategy.isTopicWriteNeeded()) { - LOGGER.info( - "Truncating topic {} after child controllers tried to roll forward to not block new versions", - kafkaTopic); - truncateKafkaTopic(kafkaTopic); - } // Verify that all regions are serving the future version after roll forward // before marking status as ONLINE @@ -3592,21 +3359,15 @@ private void truncateTopicsOptionally( boolean isTargetRegionPushWithDeferredSwap = isDeferredVersionSwap && !StringUtils.isEmpty(version.getTargetSwapRegion()); - ConcurrentPushDetectionStrategy concurrentPushDetectionStrategy = - getMultiClusterConfigs().getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); if ((failedBatchPush || nonIncPushBatchSuccess && !isDeferredVersionSwap || incPushEnabledBatchPushSuccess || isTargetRegionPushWithDeferredSwap) && !getMultiClusterConfigs().getCommonConfig().disableParentTopicTruncationUponCompletion()) { - if (concurrentPushDetectionStrategy.isTopicWriteNeeded()) { - LOGGER.info("Truncating kafka topic: {} with job status: {}", kafkaTopic, currentReturnStatus); - truncateKafkaTopic(kafkaTopic); - } if (version != null && version.getPushType().isStreamReprocessing()) { String streamReprocessingTopic = Version.composeStreamReprocessingTopic(store.getName(), version.getNumber()); LOGGER.info("Truncating kafka topic: {} with job status: {}", streamReprocessingTopic, currentReturnStatus); truncateKafkaTopic(streamReprocessingTopic); + currentReturnStatusDetails.append("Stream reprocessing topic truncated"); } - currentReturnStatusDetails.append("Parent Kafka topic truncated"); } } } @@ -3962,13 +3723,6 @@ public void killOfflinePush(String clusterName, String kafkaTopic, boolean isFor * The reason is that every errored push will call this function. */ if (maxErroredTopicNumToKeep == 0) { - // Truncate Kafka topic - LOGGER.info("Truncating topic when kill offline push job, topic: {}", kafkaTopic); - ConcurrentPushDetectionStrategy concurrentPushDetectionStrategy = - getMultiClusterConfigs().getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); - if (concurrentPushDetectionStrategy.isTopicWriteNeeded()) { - truncateKafkaTopic(kafkaTopic); - } PubSubTopic correspondingStreamReprocessingTopic = pubSubTopicRepository.getTopic(Version.composeStreamReprocessingTopicFromVersionTopic(kafkaTopic)); if (getTopicManager().containsTopic(correspondingStreamReprocessingTopic)) { diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/TopicCleanupServiceForParentController.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/TopicCleanupServiceForParentController.java deleted file mode 100644 index 10f4d120222..00000000000 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/TopicCleanupServiceForParentController.java +++ /dev/null @@ -1,82 +0,0 @@ -package com.linkedin.venice.controller.kafka; - -import com.linkedin.venice.controller.Admin; -import com.linkedin.venice.controller.VeniceControllerMultiClusterConfig; -import com.linkedin.venice.controller.stats.TopicCleanupServiceStats; -import com.linkedin.venice.exceptions.VeniceException; -import com.linkedin.venice.pubsub.PubSubClientsFactory; -import com.linkedin.venice.pubsub.PubSubTopicRepository; -import com.linkedin.venice.pubsub.api.PubSubTopic; -import com.linkedin.venice.pubsub.manager.TopicManager; -import java.util.HashMap; -import java.util.Map; -import java.util.Set; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - - -/** - * In parent controller, {@link TopicCleanupServiceForParentController} will remove all the deprecated topics: - * topic with low retention policy. - */ -public class TopicCleanupServiceForParentController extends TopicCleanupService { - private static final Logger LOGGER = LogManager.getLogger(TopicCleanupServiceForParentController.class); - private static final Map storeToCountdownForDeletion = new HashMap<>(); - - public TopicCleanupServiceForParentController( - Admin admin, - VeniceControllerMultiClusterConfig multiClusterConfigs, - PubSubTopicRepository pubSubTopicRepository, - TopicCleanupServiceStats topicCleanupServiceStats, - PubSubClientsFactory pubSubClientsFactory) { - super(admin, multiClusterConfigs, pubSubTopicRepository, topicCleanupServiceStats, pubSubClientsFactory); - } - - @Override - protected void cleanupVeniceTopics() { - Set parentFabrics = multiClusterConfigs.getParentFabrics(); - if (!parentFabrics.isEmpty()) { - for (String parentFabric: parentFabrics) { - String kafkaBootstrapServers = multiClusterConfigs.getChildDataCenterKafkaUrlMap().get(parentFabric); - cleanupVeniceTopics(getTopicManager(kafkaBootstrapServers)); - } - } else { - cleanupVeniceTopics(getTopicManager()); - } - } - - private void cleanupVeniceTopics(TopicManager topicManager) { - Map topicsWithRetention = topicManager.getAllTopicRetentions(); - Map> allStoreTopics = getAllVeniceStoreTopicsRetentions(topicsWithRetention); - allStoreTopics.forEach((storeName, topics) -> { - topics.forEach((topic, retention) -> { - if (getAdmin().isTopicTruncatedBasedOnRetention(retention)) { - // Topic may be deleted after delay - int remainingFactor = storeToCountdownForDeletion.merge( - topic.getName() + "_" + topicManager.getPubSubClusterAddress(), - delayFactor, - (oldVal, givenVal) -> oldVal - 1); - if (remainingFactor > 0) { - LOGGER.info( - "Retention policy for topic: {} is: {} ms, and it is deprecated, will delete it after {} ms.", - topic, - retention, - remainingFactor * sleepIntervalBetweenTopicListFetchMs); - } else { - LOGGER.info( - "Retention policy for topic: {} is: {} ms, and it is deprecated, will delete it now.", - topic, - retention); - storeToCountdownForDeletion.remove(topic + "_" + topicManager.getPubSubClusterAddress()); - try { - topicManager.ensureTopicIsDeletedAndBlockWithRetry(topic); - } catch (VeniceException e) { - LOGGER.warn("Caught exception when trying to delete topic: {} - {}", topic, e); // log headline of e only - // No op, will try again in the next cleanup cycle. - } - } - } - }); - }); - } -} diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/AbstractTestVeniceParentHelixAdmin.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/AbstractTestVeniceParentHelixAdmin.java index 01041de4a7d..19e376fcdb1 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/AbstractTestVeniceParentHelixAdmin.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/AbstractTestVeniceParentHelixAdmin.java @@ -26,7 +26,6 @@ import com.linkedin.venice.helix.StoragePersonaRepository; import com.linkedin.venice.helix.ZkRoutersClusterManager; import com.linkedin.venice.helix.ZkStoreConfigAccessor; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.IngestionPauseMode; import com.linkedin.venice.meta.OfflinePushStrategy; import com.linkedin.venice.meta.Store; @@ -140,10 +139,6 @@ public void setupInternalMocks() { config = mockConfig(clusterName); doReturn(1).when(config).getReplicationMetadataVersion(); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY) - .doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY) - .when(config) - .getConcurrentPushDetectionStrategy(); controllerClients .put(regionName, ControllerClient.constructClusterControllerClient(clusterName, "localhost", Optional.empty())); diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapService.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapService.java index 1887b241eb1..2fb7c7a5f17 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapService.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapService.java @@ -20,7 +20,6 @@ import com.linkedin.venice.controllerapi.UpdateStoreQueryParams; import com.linkedin.venice.exceptions.VeniceException; import com.linkedin.venice.hooks.StoreVersionLifecycleEventOutcome; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.LifecycleHooksRecord; import com.linkedin.venice.meta.LifecycleHooksRecordImpl; import com.linkedin.venice.meta.ReadWriteStoreRepository; @@ -89,8 +88,6 @@ public void setUp() { doReturn("").when(clusterConfig).getDeferredVersionSwapRegionRollforwardOrder(); doReturn(clusterConfig).when(veniceControllerMultiClusterConfig).getControllerConfig(clusterName); doReturn(1).when(clusterConfig).getDeferredVersionSwapThreadPoolSize(); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(clusterConfig) - .getConcurrentPushDetectionStrategy(); childDatacenterToUrl.put(region1, "test"); childDatacenterToUrl.put(region2, "test"); @@ -907,8 +904,6 @@ public void testMarkTargetRegionPromoted( VeniceControllerClusterConfig sequentialConfig = mock(VeniceControllerClusterConfig.class); doReturn(rolloutOrder).when(sequentialConfig).getDeferredVersionSwapRegionRollforwardOrder(); doReturn(1).when(sequentialConfig).getDeferredVersionSwapThreadPoolSize(); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(sequentialConfig) - .getConcurrentPushDetectionStrategy(); doReturn(sequentialConfig).when(veniceControllerMultiClusterConfig).getControllerConfig(clusterName); doReturn(Collections.emptyMap()).when(admin).getCurrentVersionsForMultiColos(clusterName, storeName); } @@ -995,8 +990,6 @@ public void testMarkTargetRegionPromoted_SequentialAlreadyPromoted_DoesNotCallUp VeniceControllerClusterConfig sequentialConfig = mock(VeniceControllerClusterConfig.class); doReturn(rolloutOrder).when(sequentialConfig).getDeferredVersionSwapRegionRollforwardOrder(); doReturn(1).when(sequentialConfig).getDeferredVersionSwapThreadPoolSize(); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(sequentialConfig) - .getConcurrentPushDetectionStrategy(); doReturn(sequentialConfig).when(veniceControllerMultiClusterConfig).getControllerConfig(clusterName); doReturn(Collections.emptyMap()).when(admin).getCurrentVersionsForMultiColos(clusterName, storeName); diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java index 497a0cc2d97..1a7f6a000cd 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java @@ -10,7 +10,6 @@ import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; -import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import com.linkedin.venice.controller.stats.DeferredVersionSwapStats; @@ -19,7 +18,6 @@ import com.linkedin.venice.controllerapi.StoreResponse; import com.linkedin.venice.exceptions.VeniceException; import com.linkedin.venice.hooks.StoreVersionLifecycleEventOutcome; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.LifecycleHooksRecord; import com.linkedin.venice.meta.LifecycleHooksRecordImpl; import com.linkedin.venice.meta.ReadWriteStoreRepository; @@ -419,8 +417,6 @@ public void testSequentialRolloutFailurePath() throws Exception { // Simulate failure on region2 rollout by making rollForwardToFutureVersion throw an exception when region2 appears doThrow(new VeniceException()).when(admin).rollForwardToFutureVersion(clusterName, storeName, region2); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(clusterConfig) - .getConcurrentPushDetectionStrategy(); // Create service DeferredVersionSwapService deferredVersionSwapService = new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository); @@ -599,8 +595,6 @@ public void testSequentialRolloutFinalRegionCompletion() throws Exception { String kafkaTopicName = Version.composeKafkaTopic(storeName, versionTwo); doReturn(offlinePushStatusInfoWithCompletedPush).when(admin).getOffLinePushStatus(clusterName, kafkaTopicName); - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(clusterConfig).getConcurrentPushDetectionStrategy(); - // Create service DeferredVersionSwapService deferredVersionSwapService = new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository); @@ -611,7 +605,6 @@ public void testSequentialRolloutFinalRegionCompletion() throws Exception { TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> { // Verify that updateStore was called to mark parent version as ONLINE verify(store, atLeastOnce()).updateVersionStatus(versionTwo, VersionStatus.ONLINE); - verify(admin, times(1)).truncateKafkaTopic(anyString()); }); // Verify error recording was not called diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java index 880f7d5c921..b25645aaeee 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java @@ -24,7 +24,6 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertThrows; import static org.testng.Assert.assertTrue; import static org.testng.Assert.expectThrows; @@ -54,7 +53,6 @@ import com.linkedin.venice.helix.SafeHelixDataAccessor; import com.linkedin.venice.helix.SafeHelixManager; import com.linkedin.venice.ingestion.control.RealTimeTopicSwitcher; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.HybridStoreConfig; import com.linkedin.venice.meta.Instance; import com.linkedin.venice.meta.MaterializedViewParameters; @@ -1705,64 +1703,4 @@ public void testCheckStoreGraveyardForRecreation() { assertTrue(customException.getMessage().contains("Required waiting period: 3600 seconds")); } - - @Test - public void testShouldSkipTruncatingTopicForChildControllers() { - VeniceHelixAdmin admin = mock(VeniceHelixAdmin.class); - VeniceControllerClusterConfig config = mock(VeniceControllerClusterConfig.class); - - Map configMap = new HashMap<>(); - configMap.put(clusterName, config); - doReturn(new VeniceControllerMultiClusterConfig(configMap)).when(admin).getMultiClusterConfigs(); - doReturn(false).when(admin).isParent(); - doCallRealMethod().when(admin).shouldSkipTruncatingTopic(clusterName); - - boolean shouldSkip = admin.shouldSkipTruncatingTopic(clusterName); - verify(admin, times(1)).isParent(); - assertFalse(shouldSkip); - } - - @Test - public void testShouldSkipTruncatingTopicForParentControllersTopicWriteNeeded() { - VeniceHelixAdmin admin = mock(VeniceHelixAdmin.class); - VeniceControllerClusterConfig config = mock(VeniceControllerClusterConfig.class); - Map configMap = new HashMap<>(); - configMap.put(clusterName, config); - doReturn(new VeniceControllerMultiClusterConfig(configMap)).when(admin).getMultiClusterConfigs(); - doReturn(true).when(admin).isParent(); - doCallRealMethod().when(admin).shouldSkipTruncatingTopic(clusterName); - - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); - boolean shouldSkip = admin.shouldSkipTruncatingTopic(clusterName); - verify(admin, times(1)).getMultiClusterConfigs(); - verify(admin, times(1)).isParent(); - assertFalse(shouldSkip); - - doReturn(ConcurrentPushDetectionStrategy.DUAL).when(config).getConcurrentPushDetectionStrategy(); - boolean shouldSkip2 = admin.shouldSkipTruncatingTopic(clusterName); - verify(admin, times(2)).getMultiClusterConfigs(); - verify(admin, times(2)).isParent(); - assertFalse(shouldSkip2); - - } - - @Test - public void testShouldSkipTruncatingTopicForParentControllersTopicWriteNotNeeded() { - VeniceHelixAdmin admin = mock(VeniceHelixAdmin.class); - VeniceControllerClusterConfig config = mock(VeniceControllerClusterConfig.class); - doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(config) - .getConcurrentPushDetectionStrategy(); - doReturn(true).when(admin).isParent(); - - Map configMap = new HashMap<>(); - configMap.put(clusterName, config); - doReturn(new VeniceControllerMultiClusterConfig(configMap)).when(admin).getMultiClusterConfigs(); - doCallRealMethod().when(admin).shouldSkipTruncatingTopic(clusterName); - - boolean shouldSkip = admin.shouldSkipTruncatingTopic(clusterName); - verify(admin, times(1)).getMultiClusterConfigs(); - verify(admin, times(1)).isParent(); - assertTrue(shouldSkip); - } - } diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceParentHelixAdmin.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceParentHelixAdmin.java index 97ab72d3289..4d2a400ca03 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceParentHelixAdmin.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceParentHelixAdmin.java @@ -68,7 +68,6 @@ import com.linkedin.venice.exceptions.VeniceUnsupportedOperationException; import com.linkedin.venice.helix.HelixReadWriteStoreRepository; import com.linkedin.venice.meta.BufferReplayPolicy; -import com.linkedin.venice.meta.ConcurrentPushDetectionStrategy; import com.linkedin.venice.meta.DegradedDcInfo; import com.linkedin.venice.meta.ExternalStorageReadMode; import com.linkedin.venice.meta.HybridStoreConfigImpl; @@ -724,13 +723,11 @@ public void testKillOfflinePushJob() { Store store = mock(Store.class); doReturn(store).when(internalAdmin).getStore(clusterName, pubSubTopic.getStoreName()); - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); parentAdmin.initStorageCluster(clusterName); parentAdmin.killOfflinePush(clusterName, pubSubTopic.getName(), false); verify(internalAdmin).checkPreConditionForKillOfflinePush(clusterName, pubSubTopic.getName()); - verify(internalAdmin).truncateKafkaTopic(pubSubTopic.getName()); verify(veniceWriter).put(any(), any(), anyInt(), any(), any(), anyLong(), any(), any(), any(), any()); ArgumentCaptor keyCaptor = ArgumentCaptor.forClass(byte[].class); @@ -1791,13 +1788,107 @@ public void testCompliancePushCannotKillUserPush() { -1); fail("Expected VeniceException to be thrown"); } catch (VeniceException e) { - assertTrue(e.getMessage().contains("is found and it must be terminated before another push can be started")); + assertTrue( + e.getMessage() + .contains("is still in progress and must complete or be terminated before another push can be started")); } // Verify that killOfflinePush was never called verify(mockParentAdmin, never()).killOfflinePush(clusterName, version.kafkaTopicName(), true); } + @Test + public void testDeferredVersionSwapWaitMessageIsDistinctFromConcurrentPush() { + String storeName = Utils.getUniqueString("test-store"); + VeniceParentHelixAdmin mockParentAdmin = mock(VeniceParentHelixAdmin.class); + VeniceHelixAdmin mockInternalAdmin = mock(VeniceHelixAdmin.class); + + doReturn(mockInternalAdmin).when(mockParentAdmin).getVeniceHelixAdmin(); + + Store store = new ZKStore( + storeName, + "test_owner", + 1, + PersistenceType.ROCKS_DB, + RoutingStrategy.CONSISTENT_HASH, + ReadStrategy.ANY_OF_ONLINE, + OfflinePushStrategy.WAIT_N_MINUS_ONE_REPLCIA_PER_PARTITION, + 1); + + // The existing version finished its push but is mid deferred (colo-by-colo) version swap: PUSHED on + // the parent (push complete but swap not yet done) and not yet current in every region. + String userPushId = System.currentTimeMillis() + "_https://example.com/user-job"; + VersionImpl version = new VersionImpl(storeName, 1, userPushId); + version.setVersionSwapDeferred(true); + version.setTargetSwapRegion("dc-0"); + version.setStatus(VersionStatus.PUSHED); + store.addVersion(version); + doReturn(store).when(mockParentAdmin).getStore(clusterName, storeName); + + Map currentVersionsPerRegion = new HashMap<>(); + currentVersionsPerRegion.put("dc-0", 1); + currentVersionsPerRegion.put("dc-1", 0); + doReturn(currentVersionsPerRegion).when(mockParentAdmin).getCurrentVersionsForMultiColos(clusterName, storeName); + + Map configMap = new HashMap<>(); + configMap.put(clusterName, config); + doReturn( + (LingeringStoreVersionChecker) ( + store1, + version1, + time, + controllerAdmin, + requesterCert, + identityParser) -> false).when(mockParentAdmin).getLingeringStoreVersionChecker(); + doReturn(mock(UserSystemStoreLifeCycleHelper.class)).when(mockParentAdmin).getSystemStoreLifeCycleHelper(); + doReturn(new VeniceControllerMultiClusterConfig(configMap)).when(mockParentAdmin).getMultiClusterConfigs(); + doReturn(Optional.of(version.kafkaTopicName())).when(mockParentAdmin) + .getTopicForCurrentPushJob(eq(clusterName), eq(storeName), anyBoolean(), anyBoolean()); + + // A compliance push cannot kill the existing push, so the flow reaches the rejection path. + String incomingPushId = Version.generateCompliancePushId("compliance_push"); + // Stub both the 5-arg overload and the full 17-arg impl it delegates to. + doCallRealMethod().when(mockParentAdmin).incrementVersionIdempotent(clusterName, storeName, incomingPushId, 1, 1); + doCallRealMethod().when(mockParentAdmin) + .incrementVersionIdempotent( + anyString(), + anyString(), + anyString(), + anyInt(), + anyInt(), + any(), + anyBoolean(), + anyBoolean(), + any(), + any(), + any(), + anyLong(), + any(), + anyBoolean(), + any(), + anyInt(), + anyInt()); + + HelixVeniceClusterResources mockHelixVeniceClusterResources = mock(HelixVeniceClusterResources.class); + doReturn(mockHelixVeniceClusterResources).when(mockInternalAdmin).getHelixVeniceClusterResources(clusterName); + doReturn(mock(VeniceAdminStats.class)).when(mockHelixVeniceClusterResources).getVeniceAdminStats(); + + try { + mockParentAdmin.incrementVersionIdempotent(clusterName, storeName, incomingPushId, 1, 1); + fail("Expected VeniceException to be thrown"); + } catch (VeniceException e) { + assertTrue( + e.getMessage().contains("waiting on deferred version swap"), + "Blocking message should identify the deferred version swap wait: " + e.getMessage()); + assertTrue( + e.getMessage().contains("target swap region(s): dc-0"), + "Blocking message should include the target swap region: " + e.getMessage()); + Assert.assertFalse( + e.getMessage().contains("is still in progress"), + "A deferred-swap wait must not be reported as an in-flight concurrent push: " + e.getMessage()); + } + } + @Test public void testStoreVersionCleanUpWithFewerVersions() { String storeName = "test_store"; @@ -2800,31 +2891,11 @@ private Map prepareForCurrentVersionTest(int regionCou return controllerClientMap; } - @Test - public void testGetKafkaTopicsByAge() { - String storeName = Utils.getUniqueString("test-store"); - List versionTopics = parentAdmin.getKafkaTopicsByAge(storeName); - Assert.assertTrue(versionTopics.isEmpty()); - - Set topicList = new HashSet<>(); - topicList.add(pubSubTopicRepository.getTopic(storeName + "_v1")); - topicList.add(pubSubTopicRepository.getTopic(storeName + "_v2")); - topicList.add(pubSubTopicRepository.getTopic(storeName + "_v3")); - doReturn(topicList).when(topicManager).listTopics(); - versionTopics = parentAdmin.getKafkaTopicsByAge(storeName); - Assert.assertFalse(versionTopics.isEmpty()); - PubSubTopic latestTopic = versionTopics.get(0); - assertEquals(latestTopic, pubSubTopicRepository.getTopic(storeName + "_v3")); - Assert.assertTrue(topicList.containsAll(versionTopics)); - Assert.assertTrue(versionTopics.containsAll(topicList)); - } - @Test public void testGetTopicForCurrentPushJob() { String storeName = Utils.getUniqueString("test-store"); VeniceParentHelixAdmin mockParentAdmin = mock(VeniceParentHelixAdmin.class); doReturn(internalAdmin).when(mockParentAdmin).getVeniceHelixAdmin(); - doReturn(new ArrayList()).when(mockParentAdmin).getKafkaTopicsByAge(any()); ControllerClient client = mock(ControllerClient.class); Map map = new HashMap<>(); map.put("dc-0", client); @@ -2835,8 +2906,6 @@ public void testGetTopicForCurrentPushJob() { HelixVeniceClusterResources clusterResources = internalAdmin.getHelixVeniceClusterResources(clusterName); doReturn(clusterResources).when(internalAdmin).getHelixVeniceClusterResources(clusterName); doCallRealMethod().when(mockParentAdmin).getTopicForCurrentPushJob(clusterName, storeName, false, false); - doCallRealMethod().when(mockParentAdmin) - .getTopicForCurrentPushJobParentVersionStatusBasedTracking(clusterName, storeName); Store store = new ZKStore( storeName, @@ -2848,7 +2917,9 @@ public void testGetTopicForCurrentPushJob() { OfflinePushStrategy.WAIT_N_MINUS_ONE_REPLCIA_PER_PARTITION, 1); VersionImpl version = new VersionImpl(storeName, 1, "test_push_id"); - version.setStatus(VersionStatus.ONLINE); + // STARTED is the status that polls the child job status to decide whether a push is still in + // flight; the assertions below exercise that polling/retry path. + version.setStatus(VersionStatus.STARTED); store.addVersion(version); doReturn(store).when(mockParentAdmin).getStore(clusterName, storeName); StoreResponse response = mock(StoreResponse.class); @@ -2955,79 +3026,8 @@ public void testGetTopicForCurrentPushJob() { Assert.assertFalse(mockParentAdmin.getTopicForCurrentPushJob(clusterName, storeName, false, false).isPresent()); } - @Test - public void testTruncateTopicsBasedOnMaxErroredTopicNumToKeep() { - String storeName = Utils.getUniqueString("test-store"); - VeniceParentHelixAdmin mockParentAdmin = mock(VeniceParentHelixAdmin.class); - List topics = new ArrayList<>(); - topics.add(storeName + "_v1"); - topics.add(storeName + "_v10"); - topics.add(storeName + "_v8"); - topics.add(storeName + "_v5"); - topics.add(storeName + "_v7"); - doReturn(topics).when(mockParentAdmin).existingVersionTopicsForStore(storeName); - // isTopicTruncated will return false for other topics - doReturn(true).when(mockParentAdmin).isTopicTruncated(storeName + "_v8"); - doCallRealMethod().when(mockParentAdmin).truncateTopicsBasedOnMaxErroredTopicNumToKeep(any(), anyBoolean(), any()); - doCallRealMethod().when(mockParentAdmin).setMaxErroredTopicNumToKeep(anyInt()); - mockParentAdmin.setMaxErroredTopicNumToKeep(2); - mockParentAdmin.truncateTopicsBasedOnMaxErroredTopicNumToKeep(topics, false, null); - /** - * Since the max error version topics we would like to keep is 2 and the non-truncated version - * topics include v1, v5, v7 and v10 (v8 is truncated already), we will truncate v1, v5 and keep - * 2 error non-truncated version topics v7 and v10. - */ - verify(mockParentAdmin).truncateKafkaTopic(storeName + "_v1"); - verify(mockParentAdmin).truncateKafkaTopic(storeName + "_v5"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName + "_v7"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName + "_v8"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName + "_v10"); - - // Test with more truncated topics - String storeName1 = Utils.getUniqueString("test-store"); - List topics1 = new ArrayList<>(); - topics1.add(storeName1 + "_v1"); - topics1.add(storeName1 + "_v10"); - topics1.add(storeName1 + "_v8"); - topics1.add(storeName1 + "_v5"); - topics1.add(storeName1 + "_v7"); - doReturn(topics1).when(mockParentAdmin).existingVersionTopicsForStore(storeName1); - doReturn(true).when(mockParentAdmin).isTopicTruncated(storeName1 + "_v10"); - doReturn(true).when(mockParentAdmin).isTopicTruncated(storeName1 + "_v7"); - doReturn(true).when(mockParentAdmin).isTopicTruncated(storeName1 + "_v8"); - doCallRealMethod().when(mockParentAdmin).truncateTopicsBasedOnMaxErroredTopicNumToKeep(any(), anyBoolean(), any()); - mockParentAdmin.truncateTopicsBasedOnMaxErroredTopicNumToKeep(topics1, false, null); - /** - * Since the max error version topics we would like to keep is 2 and we only have 2 non-truncated version - * topics v1 and v5 (v7, v8 and v10 are truncated already), we will not truncate anything. - */ - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName1 + "_v1"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName1 + "_v5"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName1 + "_v7"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName1 + "_v8"); - verify(mockParentAdmin, never()).truncateKafkaTopic(storeName1 + "_v10"); - } - - @Test - public void testAdminCanCleanupLeakingTopics() { - String storeName = "test_store"; - - List pubSubTopics = Arrays.asList( - pubSubTopicRepository.getTopic(storeName + "_v1"), - pubSubTopicRepository.getTopic(storeName + "_v2"), - pubSubTopicRepository.getTopic(storeName + "_v3")); - List topics = Arrays.asList(storeName + "_v1", storeName + "_v2", storeName + "_v3"); - doReturn(new HashSet(pubSubTopics)).when(topicManager).listTopics(); - - parentAdmin.truncateTopicsBasedOnMaxErroredTopicNumToKeep(topics, false, null); - verify(internalAdmin).truncateKafkaTopic(storeName + "_v1"); - verify(internalAdmin).truncateKafkaTopic(storeName + "_v2"); - verify(internalAdmin).truncateKafkaTopic(storeName + "_v3"); - } - @Test(dataProvider = "True-and-False", dataProviderClass = DataProviderUtils.class) public void testAdminCanKillLingeringVersion(boolean isIncrementalPush) { - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); try (PartialMockVeniceParentHelixAdmin partialMockParentAdmin = new PartialMockVeniceParentHelixAdmin(internalAdmin, config)) { long startTime = System.currentTimeMillis(); @@ -3038,6 +3038,7 @@ public void testAdminCanKillLingeringVersion(boolean isIncrementalPush) { String existingTopicName = storeName + "_v1"; Store store = mock(Store.class); Version version = new VersionImpl(storeName, 1, "test-push"); + version.setStatus(VersionStatus.STARTED); partialMockParentAdmin.setOfflineJobStatus(ExecutionStatus.STARTED); String newPushJobId = "new-test-push"; Version newVersion = new VersionImpl(storeName, 2, newPushJobId); @@ -3045,6 +3046,7 @@ public void testAdminCanKillLingeringVersion(boolean isIncrementalPush) { doReturn(24).when(store).getBootstrapToOnlineTimeoutInHours(); doReturn(-1).when(store).getRmdVersion(); doReturn(store).when(internalAdmin).getStore(clusterName, storeName); + doReturn(1).when(store).getLargestUsedVersionNumber(); doReturn(version).when(store).getVersion(1); doReturn(new StoreVersionInfo(store, version)).when(internalAdmin) .waitVersion(eq(clusterName), eq(storeName), eq(version.getNumber()), any()); @@ -3443,8 +3445,6 @@ public void testRollForwardSuccess() { doNothing().when(adminSpy) .sendAdminMessageAndWaitForConsumed(eq(clusterName), eq(storeName), any(AdminOperation.class)); - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); - doReturn(true).when(adminSpy).truncateKafkaTopic(Version.composeKafkaTopic(storeName, 5)); Map after = Collections.singletonMap("r1", 5); doReturn(after).when(adminSpy).getCurrentVersionsForMultiColos(clusterName, storeName); @@ -3459,7 +3459,8 @@ public void testRollForwardSuccess() { } adminSpy.rollForwardToFutureVersion(clusterName, storeName, "r1"); - verify(adminSpy).truncateKafkaTopic(Version.composeKafkaTopic(storeName, 5)); + verify(store).updateVersionStatus(5, VersionStatus.ONLINE); + verify(store).setCurrentVersion(5); } @Test(expectedExceptions = VeniceException.class, expectedExceptionsMessageRegExp = "Roll forward failed in the following regions.*") @@ -3476,9 +3477,6 @@ public void testRollForwardPartialFailure() { doNothing().when(adminSpy) .sendAdminMessageAndWaitForConsumed(eq(clusterName), eq(storeName), any(AdminOperation.class)); - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); - doReturn(true).when(adminSpy).truncateKafkaTopic(anyString()); - for (Map.Entry entry: controllerClients.entrySet()) { ControllerResponse response = new ControllerResponse(); response.setError("test error"); @@ -3500,8 +3498,6 @@ public void testRollForwardNotAllRegionsServingFutureVersionSkipsParentUpdate() doNothing().when(adminSpy) .sendAdminMessageAndWaitForConsumed(eq(clusterName), eq(storeName), any(AdminOperation.class)); - doReturn(ConcurrentPushDetectionStrategy.TOPIC_BASED_ONLY).when(config).getConcurrentPushDetectionStrategy(); - doReturn(true).when(adminSpy).truncateKafkaTopic(anyString()); // r1 rolled forward to version 5, but r2 is still on version 4 Map currentVersions = new HashMap<>(); diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForMultiKafkaClusters.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForMultiKafkaClusters.java deleted file mode 100644 index f4f3ce5ae35..00000000000 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForMultiKafkaClusters.java +++ /dev/null @@ -1,141 +0,0 @@ -package com.linkedin.venice.controller.kafka; - -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; -import static org.mockito.Mockito.verify; - -import com.linkedin.venice.controller.Admin; -import com.linkedin.venice.controller.VeniceControllerClusterConfig; -import com.linkedin.venice.controller.VeniceControllerMultiClusterConfig; -import com.linkedin.venice.controller.stats.TopicCleanupServiceStats; -import com.linkedin.venice.pubsub.PubSubClientsFactory; -import com.linkedin.venice.pubsub.PubSubTopicRepository; -import com.linkedin.venice.pubsub.adapter.kafka.admin.ApacheKafkaAdminAdapterFactory; -import com.linkedin.venice.pubsub.api.PubSubTopic; -import com.linkedin.venice.pubsub.manager.TopicManager; -import com.linkedin.venice.utils.Utils; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.ExecutionException; -import org.testng.annotations.BeforeTest; -import org.testng.annotations.Test; - - -public class TestTopicCleanupServiceForMultiKafkaClusters { - private Admin admin; - private TopicManager topicManager1; - private TopicManager topicManager2; - private TopicCleanupServiceForParentController topicCleanupService; - - private PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository(); - private final PubSubClientsFactory pubSubClientsFactory = mock(PubSubClientsFactory.class); - - @BeforeTest - public void setUp() { - VeniceControllerMultiClusterConfig config = mock(VeniceControllerMultiClusterConfig.class); - doReturn(1000l).when(config).getTopicCleanupSleepIntervalBetweenTopicListFetchMs(); - doReturn(2).when(config).getTopicCleanupDelayFactor(); - VeniceControllerClusterConfig controllerConfig = mock(VeniceControllerClusterConfig.class); - doReturn(controllerConfig).when(config).getCommonConfig(); - doReturn(new ApacheKafkaAdminAdapterFactory()).when(config).getSourceOfTruthAdminAdapterFactory(); - doReturn(Utils.setOf("fabric1", "fabric2")).when(controllerConfig).getChildDatacenters(); - - String kafkaClusterKey1 = "fabric1"; - String kafkaClusterKey2 = "fabric2"; - String kafkaClusterServerUrl1 = "host1"; - String kafkaClusterServerUrl2 = "host2"; - String kafkaClusterServerZk1 = "zk1"; - String kafkaClusterServerZk2 = "zk2"; - Set parentFabrics = new HashSet<>(); - parentFabrics.add(kafkaClusterKey1); - parentFabrics.add(kafkaClusterKey2); - Map kafkaUrlMap = new HashMap<>(); - kafkaUrlMap.put(kafkaClusterKey1, kafkaClusterServerUrl1); - kafkaUrlMap.put(kafkaClusterKey2, kafkaClusterServerUrl2); - doReturn(parentFabrics).when(config).getParentFabrics(); - doReturn(kafkaUrlMap).when(config).getChildDataCenterKafkaUrlMap(); - - admin = mock(Admin.class); - doReturn(true).when(admin).isParent(); - topicManager1 = mock(TopicManager.class); - doReturn(kafkaClusterServerUrl1).when(topicManager1).getPubSubClusterAddress(); - doReturn(topicManager1).when(admin).getTopicManager(kafkaClusterServerUrl1); - topicManager2 = mock(TopicManager.class); - doReturn(kafkaClusterServerUrl2).when(topicManager2).getPubSubClusterAddress(); - doReturn(topicManager2).when(admin).getTopicManager(kafkaClusterServerUrl2); - TopicCleanupServiceStats topicCleanupServiceStats = mock(TopicCleanupServiceStats.class); - doReturn(new ApacheKafkaAdminAdapterFactory()).when(pubSubClientsFactory).getAdminAdapterFactory(); - topicCleanupService = new TopicCleanupServiceForParentController( - admin, - config, - pubSubTopicRepository, - topicCleanupServiceStats, - pubSubClientsFactory); - } - - @Test - public void testCleanupVeniceTopics() throws ExecutionException { - Map storeTopics = new HashMap<>(); - storeTopics.put(pubSubTopicRepository.getTopic("store1_v1"), 1000l); - storeTopics.put(pubSubTopicRepository.getTopic("store1_v2"), 1000l); - storeTopics.put(pubSubTopicRepository.getTopic("store1_v3"), Long.MAX_VALUE); - storeTopics.put(pubSubTopicRepository.getTopic("store1_rt"), 1000l); - // storeTopics.put("non_venice_topic1", 1000l); - - doReturn(storeTopics).when(topicManager1).getAllTopicRetentions(); - doReturn(storeTopics).when(topicManager2).getAllTopicRetentions(); - doReturn(false).when(admin).isTopicTruncatedBasedOnRetention(Long.MAX_VALUE); - doReturn(true).when(admin).isTopicTruncatedBasedOnRetention(1000l); - - /** - * Truncated topics in parent fabrics will not be deleted in the first 2 iterations. - */ - topicCleanupService.cleanupVeniceTopics(); - - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_rt", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v1", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v2", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v3", false); - - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_rt", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v1", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v2", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v3", false); - - topicCleanupService.cleanupVeniceTopics(); - - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_rt", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v1", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v2", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v3", false); - - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_rt", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v1", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v2", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v3", false); - - topicCleanupService.cleanupVeniceTopics(); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_rt", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v1", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v2", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager1, "store1_v3", false); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_rt", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v1", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v2", true); - verifyEnsureTopicIsDeletedAndBlockWithRetry(topicManager2, "store1_v3", false); - } - - private void verifyEnsureTopicIsDeletedAndBlockWithRetry( - TopicManager topicManager, - String topicName, - boolean happened) throws ExecutionException { - if (happened) { - verify(topicManager).ensureTopicIsDeletedAndBlockWithRetry(pubSubTopicRepository.getTopic(topicName)); - } else { - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(pubSubTopicRepository.getTopic(topicName)); - } - } -} diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForParentController.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForParentController.java deleted file mode 100644 index 4893b51587e..00000000000 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/TestTopicCleanupServiceForParentController.java +++ /dev/null @@ -1,93 +0,0 @@ -package com.linkedin.venice.controller.kafka; - -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.never; -import static org.mockito.Mockito.verify; - -import com.linkedin.venice.controller.Admin; -import com.linkedin.venice.controller.VeniceControllerClusterConfig; -import com.linkedin.venice.controller.VeniceControllerMultiClusterConfig; -import com.linkedin.venice.controller.stats.TopicCleanupServiceStats; -import com.linkedin.venice.pubsub.PubSubClientsFactory; -import com.linkedin.venice.pubsub.PubSubTopicRepository; -import com.linkedin.venice.pubsub.adapter.kafka.admin.ApacheKafkaAdminAdapterFactory; -import com.linkedin.venice.pubsub.api.PubSubTopic; -import com.linkedin.venice.pubsub.manager.TopicManager; -import java.util.Collections; -import java.util.HashMap; -import java.util.Map; -import java.util.concurrent.ExecutionException; -import org.testng.annotations.BeforeTest; -import org.testng.annotations.Test; - - -public class TestTopicCleanupServiceForParentController { - private Admin admin; - private TopicManager topicManager; - private TopicCleanupService topicCleanupService; - private final PubSubTopicRepository pubSubTopicRepository = new PubSubTopicRepository(); - private final PubSubClientsFactory pubSubClientsFactory = mock(PubSubClientsFactory.class); - - @BeforeTest - public void setUp() { - admin = mock(Admin.class); - doReturn(true).when(admin).isParent(); - topicManager = mock(TopicManager.class); - doReturn(topicManager).when(admin).getTopicManager(); - VeniceControllerMultiClusterConfig config = mock(VeniceControllerMultiClusterConfig.class); - doReturn(1000l).when(config).getTopicCleanupSleepIntervalBetweenTopicListFetchMs(); - doReturn(2).when(config).getTopicCleanupDelayFactor(); - VeniceControllerClusterConfig controllerConfig = mock(VeniceControllerClusterConfig.class); - doReturn(controllerConfig).when(config).getCommonConfig(); - doReturn(Collections.singleton("dc1")).when(controllerConfig).getChildDatacenters(); - TopicCleanupServiceStats topicCleanupServiceStats = mock(TopicCleanupServiceStats.class); - doReturn(new ApacheKafkaAdminAdapterFactory()).when(pubSubClientsFactory).getAdminAdapterFactory(); - doReturn(new ApacheKafkaAdminAdapterFactory()).when(config).getSourceOfTruthAdminAdapterFactory(); - topicCleanupService = new TopicCleanupServiceForParentController( - admin, - config, - pubSubTopicRepository, - topicCleanupServiceStats, - pubSubClientsFactory); - } - - @Test - public void testCleanupVeniceTopics() throws ExecutionException { - Map storeTopics = new HashMap<>(); - PubSubTopic store1V1 = pubSubTopicRepository.getTopic("store1_v1"); - PubSubTopic store1V2 = pubSubTopicRepository.getTopic("store1_v2"); - PubSubTopic store1V3 = pubSubTopicRepository.getTopic("store1_v3"); - PubSubTopic store1RT = pubSubTopicRepository.getTopic("store1_rt"); - PubSubTopic nonVeniceTopic1 = pubSubTopicRepository.getTopic("non_venice_topic1_v1"); - storeTopics.put(store1V1, 1000l); - storeTopics.put(store1V2, 1000l); - storeTopics.put(store1V3, Long.MAX_VALUE); - storeTopics.put(store1RT, 1000l); - - doReturn(storeTopics).when(topicManager).getAllTopicRetentions(); - doReturn(false).when(admin).isTopicTruncatedBasedOnRetention(Long.MAX_VALUE); - doReturn(true).when(admin).isTopicTruncatedBasedOnRetention(1000l); - - topicCleanupService.cleanupVeniceTopics(); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1RT); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V1); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V2); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V3); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(nonVeniceTopic1); - - topicCleanupService.cleanupVeniceTopics(); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1RT); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V1); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V2); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V3); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(nonVeniceTopic1); - - topicCleanupService.cleanupVeniceTopics(); - verify(topicManager).ensureTopicIsDeletedAndBlockWithRetry(store1RT); - verify(topicManager).ensureTopicIsDeletedAndBlockWithRetry(store1V1); - verify(topicManager).ensureTopicIsDeletedAndBlockWithRetry(store1V2); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(store1V3); - verify(topicManager, never()).ensureTopicIsDeletedAndBlockWithRetry(nonVeniceTopic1); - } -}