Skip to content

Commit 29a74c8

Browse files
authored
Refactor job type declarations to use generics for improved type safety (apache#37173)
1 parent b06b9ba commit 29a74c8

17 files changed

Lines changed: 33 additions & 33 deletions

File tree

features/sharding/core/src/test/java/org/apache/shardingsphere/sharding/route/engine/condition/generator/impl/ConditionValueInOperatorGeneratorTest.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ class ConditionValueInOperatorGeneratorTest {
5252

5353
private final TimestampServiceRule timestampServiceRule = new TimestampServiceRule(new TimestampServiceRuleConfiguration("System", new Properties()));
5454

55+
@SuppressWarnings("UseOfObsoleteDateTimeApi")
5556
@Test
5657
void assertNowExpression() {
5758
ListExpression listExpression = new ListExpression(0, 0);
@@ -71,7 +72,7 @@ void assertNullExpression() {
7172
InExpression inExpression = new InExpression(0, 0, null, listExpression, false);
7273
Optional<ShardingConditionValue> shardingConditionValue = generator.generate(inExpression, column, new LinkedList<>(), timestampServiceRule);
7374
assertTrue(shardingConditionValue.isPresent());
74-
assertThat(((ListShardingConditionValue) shardingConditionValue.get()).getValues(), is(Arrays.asList(null, null)));
75+
assertThat(((ListShardingConditionValue<?>) shardingConditionValue.get()).getValues(), is(Arrays.asList(null, null)));
7576
assertTrue(shardingConditionValue.get().getParameterMarkerIndexes().isEmpty());
7677
assertThat(shardingConditionValue.get().toString(), is("tbl.id in (,)"));
7778
}
@@ -86,7 +87,7 @@ void assertNullAndCommonExpression() {
8687
InExpression inExpression = new InExpression(0, 0, null, listExpression, false);
8788
Optional<ShardingConditionValue> shardingConditionValue = generator.generate(inExpression, column, new LinkedList<>(), timestampServiceRule);
8889
assertTrue(shardingConditionValue.isPresent());
89-
assertThat(((ListShardingConditionValue) shardingConditionValue.get()).getValues(), is(Arrays.asList("test1", null, null, "test2")));
90+
assertThat(((ListShardingConditionValue<?>) shardingConditionValue.get()).getValues(), is(Arrays.asList("test1", null, null, "test2")));
9091
assertTrue(shardingConditionValue.get().getParameterMarkerIndexes().isEmpty());
9192
assertThat(shardingConditionValue.get().toString(), is("tbl.id in (test1,,,test2)"));
9293
}

kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/job/executor/DistributedPipelineJobExecutor.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ public void execute(final ShardingContext shardingContext) {
6868
log.info("Job is stopping, ignore.");
6969
return;
7070
}
71-
PipelineJobType jobType = PipelineJobIdUtils.parseJobType(jobId);
71+
PipelineJobType<?> jobType = PipelineJobIdUtils.parseJobType(jobId);
7272
PipelineContextKey contextKey = PipelineJobIdUtils.parseContextKey(jobId);
7373
PipelineJobConfiguration jobConfig = jobType.getOption().getYamlJobConfigurationSwapper().swapToObject(shardingContext.getJobParameter());
7474
PipelineJobItemManager<PipelineJobItemProgress> jobItemManager = new PipelineJobItemManager<>(jobType.getOption().getYamlJobItemProgressSwapper());
@@ -111,7 +111,7 @@ private boolean execute(final PipelineJobItemContext jobItemContext, final Pipel
111111
return true;
112112
}
113113

114-
private TransmissionProcessContext createTransmissionProcessContext(final String jobId, final PipelineJobType jobType, final PipelineContextKey contextKey) {
114+
private TransmissionProcessContext createTransmissionProcessContext(final String jobId, final PipelineJobType<?> jobType, final PipelineContextKey contextKey) {
115115
if (!jobType.getOption().isTransmissionJob()) {
116116
return null;
117117
}

kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/job/id/PipelineJobId.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ public interface PipelineJobId {
3636
*
3737
* @return pipeline job type
3838
*/
39-
PipelineJobType getJobType();
39+
PipelineJobType<?> getJobType();
4040

4141
/**
4242
* Get pipeline context key.

kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/listener/PipelineContextManagerLifecycleListener.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,7 @@ private void dispatchEnablePipelineJobStartEvent(final PipelineContextKey contex
7171
.stream().filter(each -> !each.getJobName().startsWith("_")).collect(Collectors.toList());
7272
log.info("All job names: {}", allJobsBriefInfo.stream().map(JobBriefInfo::getJobName).collect(Collectors.joining(",")));
7373
for (JobBriefInfo each : allJobsBriefInfo) {
74-
PipelineJobType jobType;
74+
PipelineJobType<?> jobType;
7575
try {
7676
jobType = PipelineJobIdUtils.parseJobType(each.getJobName());
7777
} catch (final IllegalArgumentException ex) {

kernel/data-pipeline/core/src/main/java/org/apache/shardingsphere/data/pipeline/core/task/runner/TransmissionTasksRunner.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ public final class TransmissionTasksRunner implements PipelineTasksRunner {
5454

5555
private final Collection<PipelineTask> incrementalTasks;
5656

57-
private final PipelineJobType jobType;
57+
private final PipelineJobType<?> jobType;
5858

5959
private final PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager;
6060

kernel/data-pipeline/distsql/handler/src/test/java/org/apache/shardingsphere/data/pipeline/distsql/handler/transmission/update/AlterTransmissionRuleExecutorTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ void setUp() throws ReflectiveOperationException {
7272

7373
@Test
7474
void assertExecuteUpdate() {
75-
PipelineJobType jobType = mock(PipelineJobType.class);
75+
PipelineJobType<?> jobType = mock(PipelineJobType.class);
7676
when(jobType.getType()).thenReturn(JOB_TYPE);
7777
TransmissionRuleSegment segment = new TransmissionRuleSegment();
7878
segment.setReadSegment(new ReadOrWriteSegment(5, 1000, 200, new AlgorithmSegment("READ_LIMITER", PropertiesBuilder.build(new Property("qps", "50")))));
@@ -102,7 +102,7 @@ void assertExecuteUpdate() {
102102

103103
@Test
104104
void assertExecuteUpdatePersistWhenStreamChannelIsNull() {
105-
PipelineJobType jobType = mock(PipelineJobType.class);
105+
PipelineJobType<?> jobType = mock(PipelineJobType.class);
106106
when(jobType.getType()).thenReturn(JOB_TYPE);
107107
AlterTransmissionRuleStatement sqlStatement = new AlterTransmissionRuleStatement(JOB_TYPE, new TransmissionRuleSegment());
108108
try (MockedStatic<TypedSPILoader> mockedStatic = mockStatic(TypedSPILoader.class)) {

kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/CDCJob.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ public CDCJob(final PipelineSink sink) {
101101
public void execute(final ShardingContext shardingContext) {
102102
String jobId = shardingContext.getJobName();
103103
log.info("Execute job {}", jobId);
104-
PipelineJobType jobType = PipelineJobIdUtils.parseJobType(jobId);
104+
PipelineJobType<?> jobType = PipelineJobIdUtils.parseJobType(jobId);
105105
PipelineContextKey contextKey = PipelineJobIdUtils.parseContextKey(jobId);
106106
CDCJobConfiguration jobConfig = (CDCJobConfiguration) jobType.getOption().getYamlJobConfigurationSwapper().swapToObject(shardingContext.getJobParameter());
107107
PipelineJobItemManager<TransmissionJobItemProgress> jobItemManager = new PipelineJobItemManager<>(jobType.getOption().getYamlJobItemProgressSwapper());

kernel/data-pipeline/scenario/cdc/core/src/main/java/org/apache/shardingsphere/data/pipeline/cdc/CDCJobId.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,6 @@
2222
import org.apache.shardingsphere.data.pipeline.cdc.constant.CDCSinkType;
2323
import org.apache.shardingsphere.data.pipeline.core.context.PipelineContextKey;
2424
import org.apache.shardingsphere.data.pipeline.core.job.id.PipelineJobId;
25-
import org.apache.shardingsphere.data.pipeline.core.job.type.PipelineJobType;
2625

2726
import java.util.List;
2827

@@ -33,7 +32,7 @@
3332
@Getter
3433
public final class CDCJobId implements PipelineJobId {
3534

36-
private final PipelineJobType jobType = new CDCJobType();
35+
private final CDCJobType jobType = new CDCJobType();
3736

3837
private final PipelineContextKey contextKey;
3938

kernel/data-pipeline/scenario/consistency-check/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/consistencycheck/ConsistencyCheckJobId.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import lombok.Getter;
2121
import org.apache.shardingsphere.data.pipeline.core.context.PipelineContextKey;
2222
import org.apache.shardingsphere.data.pipeline.core.job.id.PipelineJobId;
23-
import org.apache.shardingsphere.data.pipeline.core.job.type.PipelineJobType;
2423
import org.apache.shardingsphere.data.pipeline.scenario.consistencycheck.util.ConsistencyCheckSequence;
2524

2625
/**
@@ -29,7 +28,7 @@
2928
@Getter
3029
public final class ConsistencyCheckJobId implements PipelineJobId {
3130

32-
private final PipelineJobType jobType = new ConsistencyCheckJobType();
31+
private final ConsistencyCheckJobType jobType = new ConsistencyCheckJobType();
3332

3433
private final PipelineContextKey contextKey;
3534

kernel/data-pipeline/scenario/consistency-check/src/main/java/org/apache/shardingsphere/data/pipeline/scenario/consistencycheck/task/ConsistencyCheckTasksRunner.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,7 @@
5555
@Slf4j
5656
public final class ConsistencyCheckTasksRunner implements PipelineTasksRunner {
5757

58-
private final PipelineJobType jobType = new ConsistencyCheckJobType();
58+
private final ConsistencyCheckJobType jobType = new ConsistencyCheckJobType();
5959

6060
private final PipelineJobManager jobManager = new PipelineJobManager(jobType);
6161

0 commit comments

Comments
 (0)