This is an automated email from the ASF dual-hosted git repository.
terrymanu 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 996c854c2da Fix CDC selector parsing and job startup validation
(#39612)
996c854c2da is described below
commit 996c854c2da02ce322050a37d36c8b260b0991ac
Author: Liang Zhang <[email protected]>
AuthorDate: Thu Aug 27 11:14:04 2026 +0800
Fix CDC selector parsing and job startup validation (#39612)
* Add code implementation Skill and non-regression policy gates
* Add code implementation Skill and non-regression policy gates
* Fix CDC selector parsing and job startup validation
---
.../data/pipeline/cdc/api/CDCJobAPI.java | 41 ++++++++--
.../pipeline/cdc/handler/CDCBackendHandler.java | 1 -
.../pipeline/cdc/util/CDCSchemaTableUtils.java | 12 +--
.../data/pipeline/cdc/api/CDCJobAPITest.java | 89 +++++++++++++++++++++-
.../cdc/handler/CDCBackendHandlerTest.java | 45 +++++++----
.../pipeline/cdc/util/CDCSchemaTableUtilsTest.java | 22 ++++++
.../frontend/netty/CDCChannelInboundHandler.java | 16 ++--
.../netty/CDCChannelInboundHandlerTest.java | 49 +++++++++---
8 files changed, 226 insertions(+), 49 deletions(-)
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
index 6d251ee5631..839f8ce8724 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPI.java
@@ -138,8 +138,7 @@ public final class CDCJobAPI implements TransmissionJobAPI {
if
(governanceFacade.getJobFacade().getConfiguration().isExisted(jobConfig.getJobId()))
{
log.warn("CDC job already exists in registry center, ignore, job
id is `{}`", jobConfig.getJobId());
} else {
- checkDataSources(jobConfig);
- checkSchemaTableNames(jobConfig.getSchemaTableNames());
+ checkJobConfiguration(jobConfig);
governanceFacade.getJobFacade().getJob().create(jobConfig.getJobId(),
jobType.getOption().getJobClass());
JobConfigurationPOJO jobConfigPOJO =
jobConfigManager.convertToJobConfigurationPOJO(jobConfig);
jobConfigPOJO.setDisabled(true);
@@ -151,6 +150,11 @@ public final class CDCJobAPI implements TransmissionJobAPI
{
return jobConfig.getJobId();
}
+ private void checkJobConfiguration(final CDCJobConfiguration jobConfig) {
+ checkDataSources(jobConfig);
+ checkSchemaTableNames(jobConfig.getSchemaTableNames());
+ }
+
private void checkDataSources(final CDCJobConfiguration jobConfig) {
Map<String, Map<String, Object>> dataSources =
jobConfig.getDataSourceConfig().getRootConfig().getDataSources();
for (DataNode each : getDataNodes(jobConfig)) {
@@ -255,17 +259,40 @@ public final class CDCJobAPI implements
TransmissionJobAPI {
* @param sink sink
*/
public void start(final String jobId, final PipelineSink sink) {
+ CDCJobConfiguration jobConfig =
jobConfigManager.getJobConfiguration(jobId);
+ try {
+ checkJobConfiguration(jobConfig);
+ } catch (final PipelineInvalidParameterException ex) {
+ try {
+ PipelineJobRegistry.stop(jobId);
+ // CHECKSTYLE:OFF
+ } catch (final RuntimeException cleanupException) {
+ // CHECKSTYLE:ON
+ ex.addSuppressed(cleanupException);
+ }
+ try {
+ JobConfigurationPOJO jobConfigPOJO =
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+ if (!jobConfigPOJO.isDisabled()) {
+ disable(jobConfigPOJO);
+ }
+ // CHECKSTYLE:OFF
+ } catch (final RuntimeException cleanupException) {
+ // CHECKSTYLE:ON
+ ex.addSuppressed(cleanupException);
+ }
+ throw ex;
+ }
+ PipelineJobRegistry.stop(jobId);
CDCJob job = new CDCJob(sink);
PipelineJobRegistry.add(jobId, job);
- enable(jobId);
JobConfigurationPOJO jobConfigPOJO =
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+ enable(jobConfigPOJO);
OneOffJobBootstrap oneOffJobBootstrap = new
OneOffJobBootstrap(PipelineAPIFactory.getRegistryCenter(PipelineJobIdUtils.parseContextKey(jobId)),
job, jobConfigPOJO.toJobConfiguration());
job.getJobRunnerManager().setJobBootstrap(oneOffJobBootstrap);
oneOffJobBootstrap.execute();
}
- private void enable(final String jobId) {
- JobConfigurationPOJO jobConfigPOJO =
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+ private void enable(final JobConfigurationPOJO jobConfigPOJO) {
jobConfigPOJO.setDisabled(false);
jobConfigPOJO.getProps().setProperty("start_time_millis",
String.valueOf(System.currentTimeMillis()));
jobConfigPOJO.getProps().remove("stop_time");
@@ -279,6 +306,10 @@ public final class CDCJobAPI implements TransmissionJobAPI
{
*/
public void disable(final String jobId) {
JobConfigurationPOJO jobConfigPOJO =
PipelineJobIdUtils.getElasticJobConfigurationPOJO(jobId);
+ disable(jobConfigPOJO);
+ }
+
+ private void disable(final JobConfigurationPOJO jobConfigPOJO) {
jobConfigPOJO.setDisabled(true);
jobConfigPOJO.getProps().setProperty("stop_time",
LocalDateTime.now().format(DateTimeFormatterFactory.getDatetimeFormatter()));
PipelineAPIFactory.getJobConfigurationAPI(PipelineJobIdUtils.parseContextKey(jobConfigPOJO.getJobName())).updateJobConfiguration(jobConfigPOJO);
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
index 7167f58dd79..c9a34c11a99 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandler.java
@@ -123,7 +123,6 @@ public final class CDCBackendHandler {
public void startStreaming(final String jobId, final CDCConnectionContext
connectionContext, final Channel channel) {
CDCJobConfiguration cdcJobConfig =
jobConfigManager.getJobConfiguration(jobId);
ShardingSpherePreconditions.checkNotNull(cdcJobConfig, () -> new
PipelineJobNotFoundException(jobId));
- PipelineJobRegistry.stop(jobId);
ShardingSphereDatabase database =
PipelineContextManager.getProxyContext().getMetaDataContexts().getMetaData().getDatabase(cdcJobConfig.getDatabaseName());
jobAPI.start(jobId, new PipelineCDCSocketSink(channel, database,
cdcJobConfig.getSchemaTableNames()));
connectionContext.setJobId(jobId);
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
index ba986b053bb..771985bfb70 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtils.java
@@ -69,15 +69,13 @@ public final class CDCSchemaTableUtils {
Map<String, Set<String>> result = new HashMap<>();
for (SchemaTable each : schemaTables) {
if ("*".equals(each.getSchema())) {
- result.putAll(parseTableExpressionWithAllSchema(database,
systemSchemas, each));
+ mergeSchemaTables(result,
parseTableExpressionWithAllSchema(database, systemSchemas, each));
} else if ("*".equals(each.getTable())) {
result.putAll(parseTableExpressionWithAllTable(database,
each));
} else {
- String schemaName = each.getSchema();
- if (schemaName.isEmpty()) {
- schemaName =
dialectDatabaseMetaData.getSchemaOption().getDefaultSchema().orElse("");
- }
+ String schemaName = each.getSchema().isEmpty() ?
dialectDatabaseMetaData.getSchemaOption().getDefaultSchema().orElse("") :
each.getSchema();
ShardingSphereSchema schema = database.getSchema(new
IdentifierValue(schemaName));
+ ShardingSpherePreconditions.checkNotNull(schema, () -> new
SchemaNotFoundException(schemaName));
ShardingSpherePreconditions.checkNotNull(schema.getTable(each.getTable()), ()
-> new TableNotFoundException(each.getTable()));
result.computeIfAbsent(schema.getName(), ignored -> new
HashSet<>()).add(each.getTable());
}
@@ -85,6 +83,10 @@ public final class CDCSchemaTableUtils {
return result;
}
+ private static void mergeSchemaTables(final Map<String, Set<String>>
target, final Map<String, Set<String>> source) {
+ source.forEach((schemaName, tableNames) ->
target.computeIfAbsent(schemaName, ignored -> new
HashSet<>()).addAll(tableNames));
+ }
+
private static Map<String, Set<String>>
parseTableExpressionWithAllTables(final ShardingSphereDatabase database, final
Collection<String> systemSchemas) {
Map<String, Set<String>> result = new
HashMap<>(database.getAllSchemas().size(), 1F);
for (ShardingSphereSchema schema : database.getAllSchemas()) {
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
index f23a98eb03d..0444a0e10e5 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/api/CDCJobAPITest.java
@@ -103,6 +103,7 @@ import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.emptyString;
import static org.hamcrest.Matchers.hasSize;
import static org.hamcrest.Matchers.is;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -121,7 +122,7 @@ import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(AutoMockExtension.class)
-@StaticMockSettings({PipelineAPIFactory.class, PipelineJobIdUtils.class,
PipelineJobRegistry.class})
+@StaticMockSettings({PipelineAPIFactory.class, PipelineJobIdUtils.class})
@MockitoSettings(strictness = Strictness.LENIENT)
class CDCJobAPITest {
@@ -427,19 +428,30 @@ class CDCJobAPITest {
}
@Test
- void assertStartEnableDisableAndType() {
+ void assertStartEnableDisableAndType() throws ReflectiveOperationException
{
+ PipelineJobConfigurationManager jobConfigManager =
mockPersistedJobConfiguration(createJobConfiguration(1));
JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
jobConfigPOJO.setShardingTotalCount(1);
+ jobConfigPOJO.setDisabled(true);
+ jobConfigPOJO.getProps().setProperty("stop_time", "2026-08-26
00:00:00");
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
when(PipelineAPIFactory.getRegistryCenter(any())).thenReturn(mock(CoordinatorRegistryCenter.class));
PipelineSink sink = mock(PipelineSink.class);
- try (MockedConstruction<OneOffJobBootstrap> jobBootstrapConstruction =
mockConstruction(OneOffJobBootstrap.class)) {
+ try (
+ MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class);
+ MockedConstruction<OneOffJobBootstrap>
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
jobAPI.start("foo_job", sink);
-
assertThat(jobConfigPOJO.getProps().getProperty("start_time_millis"),
is(jobConfigPOJO.getProps().getProperty("start_time_millis")));
+ verify(jobConfigManager).getJobConfiguration("foo_job");
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.stop("foo_job"));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.add(eq("foo_job"), any()));
+ assertFalse(jobConfigPOJO.isDisabled());
+
assertNotNull(jobConfigPOJO.getProps().getProperty("start_time_millis"));
+ assertFalse(jobConfigPOJO.getProps().containsKey("stop_time"));
verify(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
jobAPI.disable("foo_job");
+ assertTrue(jobConfigPOJO.isDisabled());
assertNotNull(jobConfigPOJO.getProps().getProperty("stop_time"));
jobAPI.commit("foo_job");
jobAPI.rollback("foo_job");
@@ -448,6 +460,75 @@ class CDCJobAPITest {
}
}
+ @Test
+ void assertStartRejectsPersistedUnresolvedDataSource() throws
ReflectiveOperationException {
+
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_missing_ds"),
false));
+ JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+ jobConfigPOJO.setDisabled(false);
+
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+ JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+ try (
+ MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class);
+ MockedConstruction<OneOffJobBootstrap>
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
+ PipelineInvalidParameterException actual =
assertThrows(PipelineInvalidParameterException.class, () ->
jobAPI.start("foo_job", mock(PipelineSink.class)));
+ assertThat(actual.getMessage(), containsString("foo_missing_ds"));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.stop("foo_job"));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.add(anyString(), any()), never());
+ assertTrue(jobConfigPOJO.isDisabled());
+ verify(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
+ assertTrue(jobBootstrapConstruction.constructed().isEmpty());
+ }
+ }
+
+ @Test
+ void
assertStartRejectsPersistedDuplicateBareTableWithoutRepeatedGovernanceWrite()
throws ReflectiveOperationException {
+
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_ds"),
false, Arrays.asList("foo_schema.foo_tbl", "bar_schema.FOO_TBL")));
+ JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+ jobConfigPOJO.setDisabled(true);
+
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+ JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+ try (
+ MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class);
+ MockedConstruction<OneOffJobBootstrap>
jobBootstrapConstruction = mockConstruction(OneOffJobBootstrap.class)) {
+ PipelineInvalidParameterException actual =
assertThrows(PipelineInvalidParameterException.class, () ->
jobAPI.start("foo_job", mock(PipelineSink.class)));
+ assertThat(actual.getMessage(), is("There is invalid parameter
value. More than one schema table has the same table name `FOO_TBL`."));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.stop("foo_job"));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.add(anyString(), any()), never());
+ verify(jobConfigAPI, never()).updateJobConfiguration(any());
+ assertTrue(jobBootstrapConstruction.constructed().isEmpty());
+ }
+ }
+
+ @Test
+ void assertStartPreservesValidationExceptionWhenCleanupFails() throws
ReflectiveOperationException {
+
mockPersistedJobConfiguration(createJobConfiguration(Collections.singletonList("foo_missing_ds"),
false));
+ JobConfigurationPOJO jobConfigPOJO = createJobConfigurationPOJO();
+ jobConfigPOJO.setDisabled(false);
+
when(PipelineJobIdUtils.getElasticJobConfigurationPOJO("foo_job")).thenReturn(jobConfigPOJO);
+ JobConfigurationAPI jobConfigAPI = mock(JobConfigurationAPI.class);
+
when(PipelineAPIFactory.getJobConfigurationAPI(any())).thenReturn(jobConfigAPI);
+ doThrow(new IllegalStateException("governance
unavailable")).when(jobConfigAPI).updateJobConfiguration(jobConfigPOJO);
+ try (MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class)) {
+ jobRegistryMocked.when(() ->
PipelineJobRegistry.stop("foo_job")).thenThrow(new IllegalStateException("stop
unavailable"));
+ PipelineInvalidParameterException actual =
assertThrows(PipelineInvalidParameterException.class, () ->
jobAPI.start("foo_job", mock(PipelineSink.class)));
+ assertThat(actual.getMessage(), containsString("foo_missing_ds"));
+ assertThat(actual.getSuppressed().length, is(2));
+ assertThat(actual.getSuppressed()[0].getMessage(), is("stop
unavailable"));
+ assertThat(actual.getSuppressed()[1].getMessage(), is("governance
unavailable"));
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.stop("foo_job"));
+ assertTrue(jobConfigPOJO.isDisabled());
+ }
+ }
+
+ private PipelineJobConfigurationManager
mockPersistedJobConfiguration(final CDCJobConfiguration jobConfig) throws
ReflectiveOperationException {
+ PipelineJobConfigurationManager result =
mock(PipelineJobConfigurationManager.class);
+ when(result.getJobConfiguration("foo_job")).thenReturn(jobConfig);
+
Plugins.getMemberAccessor().set(CDCJobAPI.class.getDeclaredField("jobConfigManager"),
jobAPI, result);
+ return result;
+ }
+
private void putContext(final Map<String, StorageUnit> storageUnits) {
MetaDataContexts metaDataContexts = mock(MetaDataContexts.class,
RETURNS_DEEP_STUBS);
when(metaDataContexts.getMetaData().getDatabase("foo_db").getResourceMetaData().getStorageUnits()).thenReturn(storageUnits);
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
index 44126f1f566..16d31bf0ed9 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/handler/CDCBackendHandlerTest.java
@@ -63,6 +63,7 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
import org.mockito.internal.configuration.plugins.Plugins;
import java.util.ArrayList;
@@ -81,13 +82,14 @@ import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(AutoMockExtension.class)
-@StaticMockSettings({TypedSPILoader.class, DatabaseTypedSPILoader.class,
PipelineContextManager.class, CDCSchemaTableUtils.class,
PipelineDataNodeUtils.class, PipelineJobRegistry.class,
- PipelineJobIdUtils.class, CDCImporterManager.class})
+@StaticMockSettings({TypedSPILoader.class, DatabaseTypedSPILoader.class,
PipelineContextManager.class, CDCSchemaTableUtils.class,
PipelineDataNodeUtils.class, PipelineJobIdUtils.class,
+ CDCImporterManager.class})
class CDCBackendHandlerTest {
private CDCJobAPI jobAPI;
@@ -183,11 +185,15 @@ class CDCBackendHandlerTest {
mockProxyContext(mock(ShardingSphereDatabase.class));
Channel channel = mock(Channel.class);
CDCConnectionContext connectionContext = createConnectionContext();
- backendHandler.startStreaming("foo_job", connectionContext, channel);
- ArgumentCaptor<PipelineSink> sinkCaptor =
ArgumentCaptor.forClass(PipelineSink.class);
- verify(jobAPI).start(eq("foo_job"), sinkCaptor.capture());
- assertThat(((PipelineCDCSocketSink)
sinkCaptor.getValue()).getChannel(), is(channel));
- assertThat(connectionContext.getJobId(), is("foo_job"));
+ try (MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class)) {
+ jobRegistryMocked.when(() ->
PipelineJobRegistry.stop("foo_job")).thenThrow(new IllegalStateException("stop
should be owned by CDCJobAPI"));
+ backendHandler.startStreaming("foo_job", connectionContext,
channel);
+ ArgumentCaptor<PipelineSink> sinkCaptor =
ArgumentCaptor.forClass(PipelineSink.class);
+ verify(jobAPI).start(eq("foo_job"), sinkCaptor.capture());
+ assertThat(((PipelineCDCSocketSink)
sinkCaptor.getValue()).getChannel(), is(channel));
+ assertThat(connectionContext.getJobId(), is("foo_job"));
+ jobRegistryMocked.verifyNoInteractions();
+ }
}
@Test
@@ -203,9 +209,11 @@ class CDCBackendHandlerTest {
@Test
void assertStopStreamingWhenJobMissing() {
- when(PipelineJobRegistry.get("foo_job")).thenReturn(null);
- backendHandler.stopStreaming("foo_job",
DefaultChannelId.newInstance());
- verifyNoInteractions(jobAPI);
+ try (MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class)) {
+ jobRegistryMocked.when(() ->
PipelineJobRegistry.get("foo_job")).thenReturn(null);
+ backendHandler.stopStreaming("foo_job",
DefaultChannelId.newInstance());
+ verifyNoInteractions(jobAPI);
+ }
}
@Test
@@ -214,9 +222,11 @@ class CDCBackendHandlerTest {
Channel channel = mockChannel(DefaultChannelId.newInstance());
CDCJob job = mock(CDCJob.class);
when(job.getSink()).thenReturn(new PipelineCDCSocketSink(channel,
mock(ShardingSphereDatabase.class), Collections.emptyList()));
- when(PipelineJobRegistry.get("foo_job")).thenReturn(job);
- backendHandler.stopStreaming("foo_job", targetChannelId);
- verifyNoInteractions(jobAPI);
+ try (MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class)) {
+ jobRegistryMocked.when(() ->
PipelineJobRegistry.get("foo_job")).thenReturn(job);
+ backendHandler.stopStreaming("foo_job", targetChannelId);
+ verifyNoInteractions(jobAPI);
+ }
}
@Test
@@ -225,9 +235,12 @@ class CDCBackendHandlerTest {
Channel channel = mockChannel(targetChannelId);
CDCJob job = mock(CDCJob.class);
when(job.getSink()).thenReturn(new PipelineCDCSocketSink(channel,
mock(ShardingSphereDatabase.class), Collections.emptyList()));
- when(PipelineJobRegistry.get("foo_job")).thenReturn(job);
- backendHandler.stopStreaming("foo_job", targetChannelId);
- verify(jobAPI).disable("foo_job");
+ try (MockedStatic<PipelineJobRegistry> jobRegistryMocked =
mockStatic(PipelineJobRegistry.class)) {
+ jobRegistryMocked.when(() ->
PipelineJobRegistry.get("foo_job")).thenReturn(job);
+ backendHandler.stopStreaming("foo_job", targetChannelId);
+ jobRegistryMocked.verify(() ->
PipelineJobRegistry.stop("foo_job"));
+ verify(jobAPI).disable("foo_job");
+ }
}
@Test
diff --git
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
index 05500aeed79..e617a41f6b1 100644
---
a/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
+++
b/kernel/data-pipeline/scenario/cdc/core/src/test/java/org/apache/shardingsphere/data/pipeline/cdc/util/CDCSchemaTableUtilsTest.java
@@ -21,6 +21,7 @@ import
org.apache.shardingsphere.infra.config.props.ConfigurationProperties;
import
org.apache.shardingsphere.data.pipeline.cdc.protocol.request.StreamDataRequestBody.SchemaTable;
import org.apache.shardingsphere.database.connector.core.type.DatabaseType;
+import
org.apache.shardingsphere.infra.exception.kernel.metadata.SchemaNotFoundException;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.TableNotFoundException;
import
org.apache.shardingsphere.infra.metadata.database.ShardingSphereDatabase;
import
org.apache.shardingsphere.infra.metadata.database.schema.model.ShardingSphereSchema;
@@ -92,6 +93,19 @@ class CDCSchemaTableUtilsTest {
assertThat(actualResult, is(expectedResult));
}
+ @Test
+ void assertParseTableExpressionsMergesSchemaWildcardResults() {
+ ShardingSphereSchema publicSchema = mockSchema("public", "t_order",
"t_order_item");
+ ShardingSphereDatabase database =
+ new ShardingSphereDatabase("sharding_db", databaseType, null,
null, Collections.singleton(publicSchema), new ConfigurationProperties(new
Properties()));
+ List<SchemaTable> schemaTables = Arrays.asList(
+
SchemaTable.newBuilder().setSchema("*").setTable("t_order").build(),
+
SchemaTable.newBuilder().setSchema("*").setTable("t_order_item").build());
+ Map<String, Set<String>> actualResult =
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables);
+ Map<String, Set<String>> expectedResult =
Collections.singletonMap("public", new HashSet<>(Arrays.asList("t_order",
"t_order_item")));
+ assertThat(actualResult, is(expectedResult));
+ }
+
@Test
void assertParseTableExpressionsWithAllTablesInSchema() {
ShardingSphereSchema quotedSchema = mockSchema("CaseSchema",
"t_order", "t_order_item");
@@ -156,6 +170,14 @@ class CDCSchemaTableUtilsTest {
assertThrows(TableNotFoundException.class, () ->
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables));
}
+ @Test
+ void assertParseTableExpressionsWithMissingSchema() {
+ ShardingSphereSchema publicSchema = mockSchema("public", "t_order");
+ ShardingSphereDatabase database = new
ShardingSphereDatabase("sharding_db", databaseType, null, null,
Collections.singleton(publicSchema), new ConfigurationProperties(new
Properties()));
+ List<SchemaTable> schemaTables =
Collections.singletonList(SchemaTable.newBuilder().setSchema("missing").setTable("t_order").build());
+ assertThrows(SchemaNotFoundException.class, () ->
CDCSchemaTableUtils.parseTableExpressions(database, schemaTables));
+ }
+
@Test
void assertParseTableExpressionsWithoutSchemaWithWildcard() {
ShardingSphereSchema publicSchema = mockSchema("public", "t_order",
"t_order2");
diff --git
a/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
b/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
index 206e1a587c4..0684c49b1d3 100644
---
a/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
+++
b/proxy/frontend/core/src/main/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandler.java
@@ -48,8 +48,8 @@ import
org.apache.shardingsphere.database.exception.core.exception.connection.Ac
import
org.apache.shardingsphere.database.exception.core.exception.syntax.database.UnknownDatabaseException;
import org.apache.shardingsphere.infra.version.ShardingSphereVersion;
import org.apache.shardingsphere.infra.exception.ShardingSpherePreconditions;
+import
org.apache.shardingsphere.infra.exception.external.sql.ShardingSphereSQLException;
import
org.apache.shardingsphere.infra.exception.external.sql.sqlstate.XOpenSQLState;
-import
org.apache.shardingsphere.infra.exception.external.sql.type.kernel.category.PipelineSQLException;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.rule.MissingRequiredRuleException;
import org.apache.shardingsphere.infra.metadata.user.Grantee;
import org.apache.shardingsphere.infra.metadata.user.ShardingSphereUser;
@@ -181,7 +181,7 @@ public final class CDCChannelInboundHandler extends
ChannelInboundHandlerAdapter
try {
CDCResponse response =
backendHandler.streamData(request.getRequestId(), requestBody,
connectionContext, ctx.channel());
ctx.writeAndFlush(response);
- } catch (final PipelineSQLException ex) {
+ } catch (final ShardingSphereSQLException ex) {
throw new CDCExceptionWrapper(request.getRequestId(), ex);
}
}
@@ -205,10 +205,14 @@ public final class CDCChannelInboundHandler extends
ChannelInboundHandlerAdapter
if (requestBody.getStreamingId().isEmpty()) {
throw new CDCExceptionWrapper(request.getRequestId(), new
PipelineInvalidParameterException("Streaming id is empty"));
}
- String database =
backendHandler.getDatabaseNameByJobId(requestBody.getStreamingId());
- checkPrivileges(request.getRequestId(),
connectionContext.getCurrentUser().getGrantee(), database);
- backendHandler.startStreaming(requestBody.getStreamingId(),
connectionContext, ctx.channel());
- ctx.writeAndFlush(CDCResponseUtils.succeed(request.getRequestId()));
+ try {
+ String database =
backendHandler.getDatabaseNameByJobId(requestBody.getStreamingId());
+ checkPrivileges(request.getRequestId(),
connectionContext.getCurrentUser().getGrantee(), database);
+ backendHandler.startStreaming(requestBody.getStreamingId(),
connectionContext, ctx.channel());
+
ctx.writeAndFlush(CDCResponseUtils.succeed(request.getRequestId()));
+ } catch (final ShardingSphereSQLException ex) {
+ throw new CDCExceptionWrapper(request.getRequestId(), ex);
+ }
}
private void processStopStreamingRequest(final ChannelHandlerContext ctx,
final CDCRequest request, final CDCConnectionContext connectionContext) {
diff --git
a/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
b/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
index 72d03d42da1..49825306517 100644
---
a/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
+++
b/proxy/frontend/core/src/test/java/org/apache/shardingsphere/proxy/frontend/netty/CDCChannelInboundHandlerTest.java
@@ -53,7 +53,7 @@ import
org.apache.shardingsphere.database.exception.core.exception.syntax.databa
import org.apache.shardingsphere.infra.config.props.ConfigurationProperties;
import org.apache.shardingsphere.infra.config.props.ConfigurationPropertyKey;
import
org.apache.shardingsphere.infra.exception.external.sql.sqlstate.XOpenSQLState;
-import
org.apache.shardingsphere.infra.exception.external.sql.type.kernel.category.PipelineSQLException;
+import
org.apache.shardingsphere.infra.exception.kernel.metadata.SchemaNotFoundException;
import
org.apache.shardingsphere.infra.exception.kernel.metadata.rule.MissingRequiredRuleException;
import org.apache.shardingsphere.infra.metadata.database.rule.RuleMetaData;
import org.apache.shardingsphere.infra.metadata.user.Grantee;
@@ -83,13 +83,13 @@ import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
-import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
@@ -345,16 +345,27 @@ class CDCChannelInboundHandlerTest {
}
@Test
- void assertStreamDataRequestWrapsPipelineSQLException() {
- CDCConnectionContext connectionContext = new
CDCConnectionContext(user);
- channel.attr(CONNECTION_CONTEXT_KEY).set(connectionContext);
- StreamDataRequestBody.Builder bodyBuilder =
StreamDataRequestBody.newBuilder().setDatabase("logic_db");
-
bodyBuilder.addSourceSchemaTable(StreamDataRequestBody.SchemaTable.newBuilder().setSchema("schema").setTable("table").build());
- CDCRequest request =
CDCRequest.newBuilder().setType(Type.STREAM_DATA).setRequestId("stream-request").setStreamDataRequestBody(bodyBuilder.build()).build();
- when(backendHandler.streamData(any(), any(), any(),
any())).thenThrow(mock(PipelineSQLException.class));
- ChannelHandlerContext context = channel.pipeline().context(handler);
- assertThrows(CDCExceptionWrapper.class, () ->
handler.channelRead(context, request));
- assertThat(channel.attr(CONNECTION_CONTEXT_KEY).get(),
is(connectionContext));
+ void assertStreamDataRequestWithSchemaNotFoundException() {
+ channel.attr(CONNECTION_CONTEXT_KEY).set(new
CDCConnectionContext(user));
+ CDCRequest request = createStreamDataRequest("logic_db");
+ when(backendHandler.streamData(any(), any(), any(),
any())).thenThrow(new SchemaNotFoundException("foo_schema"));
+ channel.writeInbound(request);
+ CDCResponse response = readResponseSkippingGreeting();
+ assertThat(response.getStatus(), is(Status.FAILED));
+ assertThat(response.getRequestId(), is("stream-request"));
+ assertThat(response.getErrorCode(),
is(XOpenSQLState.NOT_FOUND.getValue()));
+ }
+
+ @Test
+ void assertStreamDataRequestWithInvalidParameterException() {
+ channel.attr(CONNECTION_CONTEXT_KEY).set(new
CDCConnectionContext(user));
+ CDCRequest request = createStreamDataRequest("logic_db");
+ when(backendHandler.streamData(any(), any(), any(),
any())).thenThrow(new PipelineInvalidParameterException("invalid job
configuration"));
+ channel.writeInbound(request);
+ CDCResponse response = readResponseSkippingGreeting();
+ assertThat(response.getStatus(), is(Status.FAILED));
+ assertThat(response.getRequestId(), is("stream-request"));
+ assertThat(response.getErrorCode(),
is(XOpenSQLState.INVALID_PARAMETER_VALUE.getValue()));
}
@Test
@@ -429,6 +440,20 @@ class CDCChannelInboundHandlerTest {
assertThat(response.getStatus(), is(Status.SUCCEED));
}
+ @Test
+ void assertStartStreamingRequestWithInvalidParameterException() {
+ channel.attr(CONNECTION_CONTEXT_KEY).set(new
CDCConnectionContext(user));
+
when(backendHandler.getDatabaseNameByJobId("job-1")).thenReturn("logic_db");
+ doThrow(new PipelineInvalidParameterException("invalid job
configuration")).when(backendHandler).startStreaming(any(), any(), any());
+ StartStreamingRequestBody body =
StartStreamingRequestBody.newBuilder().setStreamingId("job-1").build();
+ CDCRequest request =
CDCRequest.newBuilder().setType(Type.START_STREAMING).setRequestId("start-request").setStartStreamingRequestBody(body).build();
+ channel.writeInbound(request);
+ CDCResponse response = readResponseSkippingGreeting();
+ assertThat(response.getStatus(), is(Status.FAILED));
+ assertThat(response.getRequestId(), is("start-request"));
+ assertThat(response.getErrorCode(),
is(XOpenSQLState.INVALID_PARAMETER_VALUE.getValue()));
+ }
+
@Test
void assertStopStreamingRequestSucceed() {
CDCConnectionContext connectionContext = new
CDCConnectionContext(user);