Package org.apache.druid.segment
Class IndexMergerBase
java.lang.Object
org.apache.druid.segment.IndexMergerBase
- All Implemented Interfaces:
IndexMerger
- Direct Known Subclasses:
IndexMergerV10,IndexMergerV9
Shared segment building functionality for segments built by merging
IndexableAdapter-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionprotected static classprotected static class -
Field Summary
FieldsModifier and TypeFieldDescriptionprotected final SegmentWriteOutMediumFactoryprotected final IndexIOprotected final com.fasterxml.jackson.databind.ObjectMapperFields inherited from interface org.apache.druid.segment.IndexMerger
INVALID_ROW, log, SERIALIZER_UTILS, UNLIMITED_MAX_COLUMNS_TO_MERGE -
Constructor Summary
ConstructorsModifierConstructorDescriptionprotectedIndexMergerBase(com.fasterxml.jackson.databind.ObjectMapper mapper, IndexIO indexIO, SegmentWriteOutMediumFactory defaultSegmentWriteOutMediumFactory) -
Method Summary
Modifier and TypeMethodDescriptionprotected intgetIndexColumnCount(List<IndexableAdapter> indexableAdapters) protected intgetIndexColumnCount(IndexableAdapter indexableAdapter) protected List<List<IndexableAdapter>>getMergePhases(List<IndexableAdapter> indexes, int maxColumnsToMerge) protected abstract voidmakeColumn(SegmentFileBuilder segmentFileBuilder, String columnName, ColumnDescriptor serdeficator) Add a column specified by aColumnDescriptorto aSegmentFileBuilderprotected Map<String,DimensionHandler> makeDimensionHandlers(List<String> mergedDimensions, List<ColumnFormat> dimFormats) protected abstract FilemakeIndexFiles(List<IndexableAdapter> adapters, Metadata segmentMetadata, File outDir, ProgressIndicator progress, List<String> mergedDimensionsWithTime, IndexMergerBase.DimensionsSpecInspector dimensionsSpecInspector, List<String> mergedMetrics, Function<List<TransformableRowIterator>, TimeAndDimsIterator> rowMergerFn, IndexSpec indexSpec, SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) MergeIndexableAdapterand build the segment fileprotected TimeAndDimsIteratormakeMergedTimeAndDimsIterator(List<IndexableAdapter> indexes, List<String> mergedDimensionsWithTime, List<String> mergedMetrics, Function<List<TransformableRowIterator>, TimeAndDimsIterator> rowMergerFn, Map<String, DimensionHandler> handlers, List<DimensionMergerV9> mergers) protected voidmakeMetricsColumns(SegmentFileBuilder segmentFileBuilder, ProgressIndicator progress, List<String> mergedMetrics, Map<String, ColumnFormat> metricsTypes, List<GenericColumnSerializer> metWriters, IndexSpec indexSpec, String namePrefix) protected MetadatamakeProjections(SegmentFileBuilder segmentFileBuilder, List<AggregateProjectionMetadata> projections, List<IndexableAdapter> adapters, IndexSpec indexSpec, SegmentWriteOutMedium segmentWriteOutMedium, ProgressIndicator progress, File segmentBaseDir, Closer closer, Map<String, DimensionMergerV9> parentMergers, Metadata segmentMetadata) protected voidmakeTimeColumn(SegmentFileBuilder segmentFileBuilder, ProgressIndicator progress, GenericColumnSerializer timeWriter, IndexSpec indexSpec, String name) merge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, File outDir, DimensionsSpec dimensionsSpec, IndexSpec indexSpec, int maxColumnsToMerge) The indexes here must have the sameMetadata, otherwise an error would be thrown.protected Filemerge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, ProgressIndicator progress, SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) protected voidmergeFormat(List<IndexableAdapter> adapters, List<String> mergedDimensions, Map<String, ColumnFormat> metricTypes, List<ColumnFormat> dimFormats) protected IndexMergerBase.IndexMergeResultmergeIndexesAndWriteColumns(List<IndexableAdapter> adapters, ProgressIndicator progress, TimeAndDimsIterator timeAndDimsIterator, GenericColumnSerializer timeWriter, ArrayList<GenericColumnSerializer> metricWriters, List<DimensionMergerV9> mergers) Returns rowNumConversions, if fillRowNumConversions argument is truemergeQueryableIndex(List<QueryableIndex> indexes, boolean rollup, AggregatorFactory[] metricAggs, DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, IndexSpec indexSpecForIntermediatePersists, ProgressIndicator progress, SegmentWriteOutMediumFactory segmentWriteOutMediumFactory, int maxColumnsToMerge) Merge a collection ofQueryableIndex.protected FilemultiphaseMerge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, IndexSpec indexSpecForIntermediatePersists, ProgressIndicator progress, SegmentWriteOutMediumFactory segmentWriteOutMediumFactory, int maxColumnsToMerge) persist(IncrementalIndex index, org.joda.time.Interval dataInterval, File outDir, IndexSpec indexSpec, ProgressIndicator progress, SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) Persist an IncrementalIndex to disk in such a way that it can be loaded back up as aQueryableIndex.protected ArrayList<GenericColumnSerializer>setupMetricsWriters(SegmentWriteOutMedium segmentWriteOutMedium, List<String> mergedMetrics, Map<String, ColumnFormat> metricsTypes, IndexSpec indexSpec, String prefix) protected GenericColumnSerializersetupTimeWriter(SegmentWriteOutMedium segmentWriteOutMedium, IndexSpec indexSpec) protected abstract booleanShould the segment track columns which were specified in the ingestion spec, but not present in the processed data?protected voidwriteDimValuesAndSetupDimConversion(List<IndexableAdapter> indexes, ProgressIndicator progress, List<String> mergedDimensions, List<DimensionMergerV9> mergers) Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.druid.segment.IndexMerger
mergeQueryableIndex, mergeQueryableIndex, persist, persist
-
Field Details
-
mapper
protected final com.fasterxml.jackson.databind.ObjectMapper mapper -
indexIO
-
defaultSegmentWriteOutMediumFactory
-
-
Constructor Details
-
IndexMergerBase
protected IndexMergerBase(com.fasterxml.jackson.databind.ObjectMapper mapper, IndexIO indexIO, SegmentWriteOutMediumFactory defaultSegmentWriteOutMediumFactory)
-
-
Method Details
-
shouldStoreEmptyColumns
protected abstract boolean shouldStoreEmptyColumns()Should the segment track columns which were specified in the ingestion spec, but not present in the processed data? -
makeIndexFiles
protected abstract File makeIndexFiles(List<IndexableAdapter> adapters, @Nullable Metadata segmentMetadata, File outDir, ProgressIndicator progress, List<String> mergedDimensionsWithTime, IndexMergerBase.DimensionsSpecInspector dimensionsSpecInspector, List<String> mergedMetrics, Function<List<TransformableRowIterator>, TimeAndDimsIterator> rowMergerFn, IndexSpec indexSpec, @Nullable SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) throws IOExceptionMergeIndexableAdapterand build the segment file- Throws:
IOException
-
makeColumn
protected abstract void makeColumn(SegmentFileBuilder segmentFileBuilder, String columnName, ColumnDescriptor serdeficator) throws IOException Add a column specified by aColumnDescriptorto aSegmentFileBuilder- Throws:
IOException
-
persist
public File persist(IncrementalIndex index, org.joda.time.Interval dataInterval, File outDir, IndexSpec indexSpec, ProgressIndicator progress, @Nullable SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) throws IOException Description copied from interface:IndexMergerPersist an IncrementalIndex to disk in such a way that it can be loaded back up as aQueryableIndex. This is *not* thread-safe and havoc will ensue if this is called and writes are still occurring on the IncrementalIndex object.- Specified by:
persistin interfaceIndexMerger- Parameters:
index- the IncrementalIndex to persistdataInterval- the Interval that the data represents. Typically, this is the same as the interval from the correspondingSegmentId.outDir- the directory to persist the data toindexSpec- storage and compression optionsprogress- an object that will receive progress updatessegmentWriteOutMediumFactory- controls allocation of temporary data structures- Returns:
- the index output directory
- Throws:
IOException- if an IO error occurs persisting the index
-
mergeQueryableIndex
public File mergeQueryableIndex(List<QueryableIndex> indexes, boolean rollup, AggregatorFactory[] metricAggs, @Nullable DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, IndexSpec indexSpecForIntermediatePersists, ProgressIndicator progress, @Nullable SegmentWriteOutMediumFactory segmentWriteOutMediumFactory, int maxColumnsToMerge) throws IOException Description copied from interface:IndexMergerMerge a collection ofQueryableIndex.- Specified by:
mergeQueryableIndexin interfaceIndexMerger- Throws:
IOException
-
merge
public File merge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, File outDir, DimensionsSpec dimensionsSpec, IndexSpec indexSpec, int maxColumnsToMerge) throws IOException The indexes here must have the sameMetadata, otherwise an error would be thrown.- Specified by:
mergein interfaceIndexMerger- Throws:
IOException
-
multiphaseMerge
protected File multiphaseMerge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, @Nullable DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, IndexSpec indexSpecForIntermediatePersists, ProgressIndicator progress, @Nullable SegmentWriteOutMediumFactory segmentWriteOutMediumFactory, int maxColumnsToMerge) throws IOException - Throws:
IOException
-
getMergePhases
protected List<List<IndexableAdapter>> getMergePhases(List<IndexableAdapter> indexes, int maxColumnsToMerge) -
getIndexColumnCount
-
getIndexColumnCount
-
merge
protected File merge(List<IndexableAdapter> indexes, boolean rollup, AggregatorFactory[] metricAggs, @Nullable DimensionsSpec dimensionsSpec, File outDir, IndexSpec indexSpec, ProgressIndicator progress, @Nullable SegmentWriteOutMediumFactory segmentWriteOutMediumFactory) throws IOException - Throws:
IOException
-
makeProjections
protected Metadata makeProjections(SegmentFileBuilder segmentFileBuilder, List<AggregateProjectionMetadata> projections, List<IndexableAdapter> adapters, IndexSpec indexSpec, SegmentWriteOutMedium segmentWriteOutMedium, ProgressIndicator progress, File segmentBaseDir, Closer closer, Map<String, DimensionMergerV9> parentMergers, Metadata segmentMetadata) throws IOException- Throws:
IOException
-
makeTimeColumn
protected void makeTimeColumn(SegmentFileBuilder segmentFileBuilder, ProgressIndicator progress, GenericColumnSerializer timeWriter, IndexSpec indexSpec, String name) throws IOException - Throws:
IOException
-
makeMetricsColumns
protected void makeMetricsColumns(SegmentFileBuilder segmentFileBuilder, ProgressIndicator progress, List<String> mergedMetrics, Map<String, ColumnFormat> metricsTypes, List<GenericColumnSerializer> metWriters, IndexSpec indexSpec, String namePrefix) throws IOException- Throws:
IOException
-
mergeIndexesAndWriteColumns
protected IndexMergerBase.IndexMergeResult mergeIndexesAndWriteColumns(List<IndexableAdapter> adapters, ProgressIndicator progress, TimeAndDimsIterator timeAndDimsIterator, GenericColumnSerializer timeWriter, ArrayList<GenericColumnSerializer> metricWriters, List<DimensionMergerV9> mergers) throws IOException Returns rowNumConversions, if fillRowNumConversions argument is true- Throws:
IOException
-
setupTimeWriter
protected GenericColumnSerializer setupTimeWriter(SegmentWriteOutMedium segmentWriteOutMedium, IndexSpec indexSpec) throws IOException - Throws:
IOException
-
setupMetricsWriters
protected ArrayList<GenericColumnSerializer> setupMetricsWriters(SegmentWriteOutMedium segmentWriteOutMedium, List<String> mergedMetrics, Map<String, ColumnFormat> metricsTypes, IndexSpec indexSpec, String prefix) throws IOException- Throws:
IOException
-
writeDimValuesAndSetupDimConversion
protected void writeDimValuesAndSetupDimConversion(List<IndexableAdapter> indexes, ProgressIndicator progress, List<String> mergedDimensions, List<DimensionMergerV9> mergers) throws IOException - Throws:
IOException
-
mergeFormat
protected void mergeFormat(List<IndexableAdapter> adapters, List<String> mergedDimensions, Map<String, ColumnFormat> metricTypes, List<ColumnFormat> dimFormats) -
makeDimensionHandlers
protected Map<String,DimensionHandler> makeDimensionHandlers(List<String> mergedDimensions, List<ColumnFormat> dimFormats) -
makeMergedTimeAndDimsIterator
protected TimeAndDimsIterator makeMergedTimeAndDimsIterator(List<IndexableAdapter> indexes, List<String> mergedDimensionsWithTime, List<String> mergedMetrics, Function<List<TransformableRowIterator>, TimeAndDimsIterator> rowMergerFn, Map<String, DimensionHandler> handlers, List<DimensionMergerV9> mergers)
-