This is an automated email from the ASF dual-hosted git repository.
zhonghongsheng pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shardingsphere.git
The following commit(s) were added to refs/heads/master by this push:
new 309d3f5 For #14722: Release schema level scaling lock for some
failure cases (#15114)
309d3f5 is described below
commit 309d3f52026c6388915f69e02c9e15d022325b37
Author: ReyYang <[email protected]>
AuthorDate: Thu Jan 27 19:57:55 2022 +0800
For #14722: Release schema level scaling lock for some failure cases
(#15114)
---
.../scenario/rulealtered/RuleAlteredJob.java | 4 +++
.../rulealtered/RuleAlteredJobScheduler.java | 6 +++++
.../subscriber/ScalingRegistrySubscriber.java | 15 +++++++++++
.../rule/ScalingReleaseSchemaNameLockEvent.java | 30 ++++++++++++++++++++++
4 files changed, 55 insertions(+)
diff --git
a/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJob.java
b/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJob.java
index f210711..eefe634 100644
---
a/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJob.java
+++
b/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJob.java
@@ -24,7 +24,9 @@ import
org.apache.shardingsphere.data.pipeline.core.api.GovernanceRepositoryAPI;
import org.apache.shardingsphere.data.pipeline.core.api.PipelineAPIFactory;
import org.apache.shardingsphere.elasticjob.api.ShardingContext;
import org.apache.shardingsphere.elasticjob.simple.job.SimpleJob;
+import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
import org.apache.shardingsphere.infra.yaml.engine.YamlEngine;
+import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.ScalingReleaseSchemaNameLockEvent;
/**
* Rule altered job.
@@ -54,6 +56,8 @@ public final class RuleAlteredJob implements SimpleJob {
jobContext.close();
jobContext.setStatus(JobStatus.PREPARING_FAILURE);
governanceRepositoryAPI.persistJobProgress(jobContext);
+ ScalingReleaseSchemaNameLockEvent event = new
ScalingReleaseSchemaNameLockEvent(jobConfig.getWorkflowConfig().getSchemaName());
+ ShardingSphereEventBus.getInstance().post(event);
throw ex;
}
governanceRepositoryAPI.persistJobProgress(jobContext);
diff --git
a/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJobScheduler.java
b/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJobScheduler.java
index e48fc0e..f89292b 100644
---
a/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJobScheduler.java
+++
b/shardingsphere-kernel/shardingsphere-data-pipeline/shardingsphere-data-pipeline-core/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/rulealtered/RuleAlteredJobScheduler.java
@@ -25,6 +25,8 @@ import
org.apache.shardingsphere.data.pipeline.api.job.JobStatus;
import org.apache.shardingsphere.data.pipeline.core.execute.ExecuteCallback;
import org.apache.shardingsphere.data.pipeline.core.task.IncrementalTask;
import org.apache.shardingsphere.data.pipeline.core.task.InventoryTask;
+import org.apache.shardingsphere.infra.eventbus.ShardingSphereEventBus;
+import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.ScalingReleaseSchemaNameLockEvent;
/**
* Rule altered job scheduler.
@@ -110,6 +112,8 @@ public final class RuleAlteredJobScheduler implements
Runnable {
log.error("Inventory task execute failed.", throwable);
stop();
jobContext.setStatus(JobStatus.EXECUTE_INVENTORY_TASK_FAILURE);
+ ScalingReleaseSchemaNameLockEvent event = new
ScalingReleaseSchemaNameLockEvent(jobContext.getJobConfig().getWorkflowConfig().getSchemaName());
+ ShardingSphereEventBus.getInstance().post(event);
}
};
}
@@ -142,6 +146,8 @@ public final class RuleAlteredJobScheduler implements
Runnable {
log.error("Incremental task execute failed.", throwable);
stop();
jobContext.setStatus(JobStatus.EXECUTE_INCREMENTAL_TASK_FAILURE);
+ ScalingReleaseSchemaNameLockEvent event = new
ScalingReleaseSchemaNameLockEvent(jobContext.getJobConfig().getWorkflowConfig().getSchemaName());
+ ShardingSphereEventBus.getInstance().post(event);
}
};
}
diff --git
a/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/cache/subscriber/ScalingRegistrySubscriber.java
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/cache/subscriber/ScalingRegistrySubscriber.java
index d286517..9bc1f9b 100644
---
a/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/cache/subscriber/ScalingRegistrySubscriber.java
+++
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/cache/subscriber/ScalingRegistrySubscriber.java
@@ -28,6 +28,7 @@ import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.cache
import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.cache.event.StartScalingEvent;
import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.ClusterSwitchConfigurationEvent;
import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.RuleConfigurationCachedEvent;
+import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.ScalingReleaseSchemaNameLockEvent;
import
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule.ScalingTaskFinishedEvent;
import
org.apache.shardingsphere.mode.metadata.persist.service.impl.DataSourcePersistService;
import
org.apache.shardingsphere.mode.metadata.persist.service.impl.SchemaRulePersistService;
@@ -131,6 +132,20 @@ public final class ScalingRegistrySubscriber {
persistService.persist(schemaName, event.getTargetRuleConfigs());
}
+ /**
+ * scaling release schema name lock.
+ *
+ * @param event scaling release schema name lock event
+ */
+ @Subscribe
+ public void scalingReleaseSchemaNameLock(final
ScalingReleaseSchemaNameLockEvent event) {
+ if (schemaNameLockedMap.getOrDefault(event.getSchemaName(), false)) {
+ log.info("scaling job finished, release schema name lock, event =
{}", event);
+ schemaNameLockedMap.remove(event.getSchemaName());
+
lockRegistryService.releaseLock(decorateLockName(event.getSchemaName()));
+ }
+ }
+
private String decorateLockName(final String schemaName) {
return "Scaling-" + schemaName;
}
diff --git
a/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/config/event/rule/ScalingReleaseSchemaNameLockEvent.java
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/config/event/rule/ScalingReleaseSchemaNameLoc
[...]
new file mode 100644
index 0000000..4d5245d
--- /dev/null
+++
b/shardingsphere-mode/shardingsphere-mode-type/shardingsphere-cluster-mode/shardingsphere-cluster-mode-core/src/main/java/org/apache/shardingsphere/mode/manager/cluster/coordinator/registry/config/event/rule/ScalingReleaseSchemaNameLockEvent.java
@@ -0,0 +1,30 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package
org.apache.shardingsphere.mode.manager.cluster.coordinator.registry.config.event.rule;
+
+import lombok.Getter;
+import lombok.RequiredArgsConstructor;
+
+/**
+ * Scaling release schema name lock event.
+ */
+@RequiredArgsConstructor
+@Getter
+public final class ScalingReleaseSchemaNameLockEvent {
+ private final String schemaName;
+}