Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,6 @@
import org.apache.druid.segment.virtual.ExpressionVirtualColumn;
import org.apache.druid.server.compaction.InlineReindexingRuleProvider;
import org.apache.druid.server.compaction.ReindexingDeletionRule;
import org.apache.druid.server.compaction.ReindexingIOConfigRule;
import org.apache.druid.server.compaction.ReindexingSegmentGranularityRule;
import org.apache.druid.server.compaction.ReindexingTuningConfigRule;
import org.apache.druid.server.coordinator.ClusterCompactionConfig;
Expand Down Expand Up @@ -244,14 +243,13 @@ public void test_compaction_withPersistLastCompactionStateFalse_storesOnlyFinger
verifySegmentsHaveNullLastCompactionStateAndNonNullFingerprint();
}

@MethodSource("getEngine")
@ParameterizedTest(name = "compactionEngine={0}")
public void test_cascadingCompactionTemplate_multiplePeriodsApplyDifferentCompactionRules(CompactionEngine compactionEngine)
@Test
public void test_cascadingCompactionTemplate_multiplePeriodsApplyDifferentCompactionRules()
{
// Configure cluster with storeCompactionStatePerSegment=false
// Configure cluster with MSQ engine and storeCompactionStatePerSegment=false
final UpdateResponse updateResponse = cluster.callApi().onLeaderOverlord(
o -> o.updateClusterCompactionConfig(
new ClusterCompactionConfig(1.0, 100, null, true, compactionEngine, false)
new ClusterCompactionConfig(1.0, 100, null, true, CompactionEngine.MSQ, false)
)
);
Assertions.assertTrue(updateResponse.isSuccess());
Expand Down Expand Up @@ -314,23 +312,18 @@ public void test_cascadingCompactionTemplate_multiplePeriodsApplyDifferentCompac
null
);

InlineReindexingRuleProvider.Builder ruleProvider = InlineReindexingRuleProvider.builder()
.segmentGranularityRules(List.of(hourRule, dayRule))
.tuningConfigRules(List.of(tuningConfigRule))
.deletionRules(List.of(deletionRule));

if (compactionEngine == CompactionEngine.NATIVE) {
ruleProvider = ruleProvider.ioConfigRules(
List.of(new ReindexingIOConfigRule("dropExisting", null, Period.days(7), new UserCompactionTaskIOConfig(true)))
);
}
InlineReindexingRuleProvider ruleProvider = InlineReindexingRuleProvider
.builder()
.segmentGranularityRules(List.of(hourRule, dayRule))
.tuningConfigRules(List.of(tuningConfigRule))
.deletionRules(List.of(deletionRule))
.build();

CascadingReindexingTemplate cascadingReindexingTemplate = new CascadingReindexingTemplate(
dataSource,
null,
null,
ruleProvider.build(),
compactionEngine,
ruleProvider,
null,
null,
null,
Expand Down Expand Up @@ -409,7 +402,6 @@ public void test_cascadingReindexing_withVirtualColumnOnNestedData_filtersCorrec
.deletionRules(List.of(deletionRule))
.tuningConfigRules(List.of(tuningConfigRule))
.build(),
compactionEngine,
null,
null,
null,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,8 +92,6 @@ public class CascadingReindexingTemplate implements CompactionJobTemplate, DataS
private final ReindexingRuleProvider ruleProvider;
@Nullable
private final Map<String, Object> taskContext;
@Nullable
private final CompactionEngine engine;
private final int taskPriority;
private final long inputSegmentSizeBytes;
private final Period skipOffsetFromLatest;
Expand All @@ -106,7 +104,6 @@ public CascadingReindexingTemplate(
@JsonProperty("taskPriority") @Nullable Integer taskPriority,
@JsonProperty("inputSegmentSizeBytes") @Nullable Long inputSegmentSizeBytes,
@JsonProperty("ruleProvider") ReindexingRuleProvider ruleProvider,
@JsonProperty("engine") @Nullable CompactionEngine engine,
@JsonProperty("taskContext") @Nullable Map<String, Object> taskContext,
@JsonProperty("skipOffsetFromLatest") @Nullable Period skipOffsetFromLatest,
@JsonProperty("skipOffsetFromNow") @Nullable Period skipOffsetFromNow,
Expand All @@ -119,7 +116,6 @@ public CascadingReindexingTemplate(
InvalidInput.conditionalException(ruleProvider != null, "'ruleProvider' cannot be null");
this.ruleProvider = ruleProvider;

this.engine = engine;
this.taskContext = taskContext;
this.taskPriority = Objects.requireNonNullElse(taskPriority, DEFAULT_COMPACTION_TASK_PRIORITY);
this.inputSegmentSizeBytes = Objects.requireNonNullElse(inputSegmentSizeBytes, DEFAULT_INPUT_SEGMENT_SIZE_BYTES);
Expand Down Expand Up @@ -149,12 +145,10 @@ public Map<String, Object> getTaskContext()
return taskContext;
}

@JsonProperty
@Nullable
@Override
public CompactionEngine getEngine()
{
return engine;
return CompactionEngine.MSQ;
}

@JsonProperty
Expand Down Expand Up @@ -347,7 +341,7 @@ private InlineSchemaDataSourceCompactionConfig.Builder createBaseBuilder()
.forDataSource(dataSource)
.withTaskPriority(taskPriority)
.withInputSegmentSizeBytes(inputSegmentSizeBytes)
.withEngine(engine)
.withEngine(CompactionEngine.MSQ)
.withTaskContext(taskContext)
.withSkipOffsetFromLatest(Period.ZERO); // We handle skip offsets at the timeline level, we know we want to cover the entirety of the interval
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
import org.apache.druid.server.compaction.IntervalGranularityInfo;
import org.apache.druid.server.compaction.ReindexingDataSchemaRule;
import org.apache.druid.server.compaction.ReindexingDeletionRule;
import org.apache.druid.server.compaction.ReindexingIOConfigRule;
import org.apache.druid.server.compaction.ReindexingRule;
import org.apache.druid.server.compaction.ReindexingRuleProvider;
import org.apache.druid.server.compaction.ReindexingTuningConfigRule;
Expand Down Expand Up @@ -137,14 +136,6 @@ BuildResult applyToWithDetails(InlineSchemaDataSourceCompactionConfig.Builder bu
count++;
}

// Apply IO config rule
ReindexingIOConfigRule ioConfigRule = provider.getIOConfigRule(interval, referenceTime);
if (ioConfigRule != null) {
builder.withIoConfig(ioConfigRule.getIoConfig());
appliedRules.add(ioConfigRule);
count++;
}

// Apply data schema rules
ReindexingDataSchemaRule dataSchemaRule = provider.getDataSchemaRule(interval, referenceTime);
if (dataSchemaRule != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,6 @@ public void test_serde() throws Exception
)
))
.build(),
CompactionEngine.NATIVE,
ImmutableMap.of("context_key", "context_value"),
null,
null,
Expand Down Expand Up @@ -123,7 +122,6 @@ public void test_serde_asDataSourceCompactionConfig() throws Exception
)
))
.build(),
CompactionEngine.MSQ,
ImmutableMap.of("key", "value"),
null,
null,
Expand Down Expand Up @@ -161,7 +159,6 @@ public void test_createCompactionJobs_ruleProviderNotReady()
null,
null,
null,
null,
Granularities.DAY
);

Expand All @@ -186,7 +183,6 @@ public void test_constructor_setBothSkipOffsetStrategiesThrowsException()
null,
mockProvider,
null,
null,
Period.days(7), // skipOffsetFromLatest
Period.days(3), // skipOffsetFromNow
Granularities.DAY
Expand All @@ -213,7 +209,6 @@ public void test_constructor_nullDataSourceThrowsException()
null,
null,
null,
null,
Granularities.DAY
)
);
Expand All @@ -235,7 +230,6 @@ public void test_constructor_nullRuleProviderThrowsException()
null,
null,
null,
null,
Granularities.DAY
)
);
Expand All @@ -259,7 +253,6 @@ public void test_constructor_nullDefaultSegmentGranularityThrowsException()
null,
null,
null,
null,
null // null defaultSegmentGranularity
)
);
Expand All @@ -278,7 +271,7 @@ public void test_createCompactionJobs_simple()
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, null, null
"testDS", null, null, mockProvider, null, null, null
);

template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -304,7 +297,7 @@ public void test_createCompactionJobs_withSkipOffsetFromLatest_skipAllOfTime()
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, Period.days(100), null
"testDS", null, null, mockProvider, null, Period.days(100), null
);

List<CompactionJob> jobs = template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -325,7 +318,7 @@ public void test_createCompactionJobs_withSkipOffsetFromLatest_skipsIntervalsExt
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, Period.days(5), null
"testDS", null, null, mockProvider, null, Period.days(5), null
);

template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -348,7 +341,7 @@ public void test_createCompactionJobs_withSkipOffsetFromLatest_eliminatesInterva
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, Period.days(15), null
"testDS", null, null, mockProvider, null, Period.days(15), null
);

template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -371,7 +364,7 @@ public void test_createCompactionJobs_withSkipOffsetFromNow_skipAllOfTime()
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, null, Period.days(100)
"testDS", null, null, mockProvider, null, null, Period.days(100)
);

List<CompactionJob> jobs = template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -392,7 +385,7 @@ public void test_createCompactionJobs_withSkipOffsetFromNow_skipsIntervalsExtend
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, null, Period.days(20)
"testDS", null, null, mockProvider, null, null, Period.days(20)
);

template.createCompactionJobs(mockSource, mockParams);
Expand All @@ -415,7 +408,7 @@ public void test_createCompactionJobs_withSkipOffsetFromNow_eliminatesInterval()
DruidInputSource mockSource = createMockSource();

TestCascadingReindexingTemplate template = new TestCascadingReindexingTemplate(
"testDS", null, null, mockProvider, null, null, null, Period.days(20)
"testDS", null, null, mockProvider, null, null, Period.days(20)
);

template.createCompactionJobs(mockSource, mockParams);
Expand Down Expand Up @@ -481,7 +474,6 @@ public void test_generateAlignedSearchIntervals_withGranularityAlignment()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -572,7 +564,6 @@ public void test_generateAlignedSearchIntervals_withNonSegmentGranularityRuleSpl
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -672,7 +663,6 @@ public void test_generateAlignedSearchIntervals_withNoSegmentGranularityRules()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -764,7 +754,6 @@ public void test_generateAlignedSearchIntervals_prependIntervalForShortNonSegmen
null,
null,
null,
null,
Granularities.HOUR
);

Expand Down Expand Up @@ -870,7 +859,6 @@ public void test_generateAlignedSearchIntervals()
null,
null,
null,
null,
Granularities.HOUR
);

Expand Down Expand Up @@ -940,7 +928,6 @@ public void test_generateAlignedSearchIntervals_noRulesThrowsException()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1005,7 +992,6 @@ public void test_generateAlignedSearchIntervals_splitPointSnapsToExistingBoundar
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1072,7 +1058,6 @@ public void test_generateAlignedSearchIntervals_prependAlignmentDoesNotExtendTim
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1143,7 +1128,6 @@ public void test_generateAlignedSearchIntervals_duplicateSplitPointsFiltered()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1207,7 +1191,6 @@ public void test_generateAlignedSearchIntervals_singleRuleOnly()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1271,7 +1254,6 @@ public void test_generateAlignedSearchIntervals_zeroPeriodRuleAppliesImmediately
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1356,7 +1338,6 @@ public void test_generateAlignedSearchIntervals_zeroPeriodRuleWithOtherRules()
null,
null,
null,
null,
Granularities.DAY
);

Expand Down Expand Up @@ -1429,7 +1410,6 @@ public void test_generateAlignedSearchIntervals_failsWhenDefaultGranularityIsCoa
null,
null,
null,
null,
Granularities.MONTH // MONTH is coarser than HOUR!
);

Expand Down Expand Up @@ -1490,7 +1470,6 @@ public void test_generateAlignedSearchIntervals_failsWhenOlderRuleHasFinerGranul
null,
null,
null,
null,
Granularities.DAY
);

Expand All @@ -1517,14 +1496,13 @@ public TestCascadingReindexingTemplate(
Integer taskPriority,
Long inputSegmentSizeBytes,
ReindexingRuleProvider ruleProvider,
CompactionEngine engine,
Map<String, Object> taskContext,
Period skipOffsetFromLatest,
Period skipOffsetFromNow
)
{
super(dataSource, taskPriority, inputSegmentSizeBytes, ruleProvider,
engine, taskContext, skipOffsetFromLatest, skipOffsetFromNow, Granularities.DAY);
taskContext, skipOffsetFromLatest, skipOffsetFromNow, Granularities.DAY);
}

public List<Interval> getProcessedIntervals()
Expand Down Expand Up @@ -1597,7 +1575,6 @@ private ReindexingRuleProvider createMockProvider(List<Period> periods)
// Return a fresh stream on each call to avoid "stream has already been operated upon or closed" errors
EasyMock.expect(mockProvider.streamAllRules()).andAnswer(() -> segmentGranularityRules.stream().map(r -> (ReindexingRule) r)).anyTimes();
EasyMock.expect(mockProvider.getSegmentGranularityRule(EasyMock.anyObject(), EasyMock.anyObject())).andReturn(segmentGranularityRules.get(0)).anyTimes();
EasyMock.expect(mockProvider.getIOConfigRule(EasyMock.anyObject(), EasyMock.anyObject())).andReturn(null).anyTimes();
EasyMock.expect(mockProvider.getTuningConfigRule(EasyMock.anyObject(), EasyMock.anyObject())).andReturn(null).anyTimes();
EasyMock.expect(mockProvider.getDataSchemaRule(EasyMock.anyObject(), EasyMock.anyObject())).andReturn(null).anyTimes();
EasyMock.expect(mockProvider.getDeletionRules(EasyMock.anyObject(), EasyMock.anyObject())).andReturn(Collections.emptyList()).anyTimes();
Expand Down
Loading
Loading