This is an automated email from the ASF dual-hosted git repository.

mcvsubbu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 9b5775e  [Issue #4551] Remove memory allocation for virtual columns in 
consumi… (#4555)
9b5775e is described below

commit 9b5775e79c8881d06b2ed472698eabbff56317d6
Author: Subbu Subramaniam <[email protected]>
AuthorDate: Sat Aug 24 15:28:20 2019 -0700

    [Issue #4551] Remove memory allocation for virtual columns in consumi… 
(#4555)
    
    * [Issue #4551] Remove memory allocation for virtual columns in consuming 
segments
    
    Removed allocation of memory for virtual columns in consuming segments.
    
    Removed the getRecord() interface from IndexSegment since it is not used 
anywhere.
    Restricted it to a method within MutableSegmentImpl (needed to build 
completed
    segments from consuming segments).
    
    Verified that segment metadata did not have virtual columns before and 
after my change.
    
    Verified that we can still execute the following queries:
    
            SELECT distinctcount($segmentName) from mytable
            SELECT count(*) from mytable group by $segmentName top 50
            SELECT somefield, $docId, $segmentName, $hostName from mytable
            SELECT count(*) from mytable where $segmentName = 
"mytable__1__23__20190823T0142Z"
    
    and they work fine
    
    * Addressed review comments
    
    Also added more test cases to cover the case of time column, and virtual
    columns being handled right in metrics aggregation
    
    * Addressed review comments
    
    * Addressed review comments
    
    * Keeping track of fieldspec for metrics and dimension columns as well.
    
    Avoid duplicate work while indexing rows.
---
 .../pinot/core/indexsegment/IndexSegment.java      |  2 +-
 .../immutable/ImmutableSegmentImpl.java            | 14 ++--
 .../indexsegment/mutable/MutableSegmentImpl.java   | 86 +++++++++++++++-------
 .../MutableSegmentImplAggregateMetricsTest.java    |  8 +-
 4 files changed, 76 insertions(+), 34 deletions(-)

diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/IndexSegment.java 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/IndexSegment.java
index 692b017..a23b43c 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/IndexSegment.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/IndexSegment.java
@@ -71,7 +71,7 @@ public interface IndexSegment {
   List<StarTreeV2> getStarTrees();
 
   /**
-   * Returns the record for the given document Id.
+   * Returns the record for the given document Id. Virtual column values are 
not returned.
    * <p>NOTE: don't use this method for high performance code.
    *
    * @param docId Document Id
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/immutable/ImmutableSegmentImpl.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/immutable/ImmutableSegmentImpl.java
index 40aa372..04c55fd 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/immutable/ImmutableSegmentImpl.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/immutable/ImmutableSegmentImpl.java
@@ -25,6 +25,7 @@ import java.util.Map;
 import java.util.Set;
 import javax.annotation.Nullable;
 import org.apache.pinot.common.data.FieldSpec;
+import org.apache.pinot.common.data.Schema;
 import org.apache.pinot.core.data.GenericRow;
 import org.apache.pinot.core.indexsegment.IndexSegmentUtils;
 import org.apache.pinot.core.io.reader.DataFileReader;
@@ -161,12 +162,15 @@ public class ImmutableSegmentImpl implements 
ImmutableSegment {
 
   @Override
   public GenericRow getRecord(int docId, GenericRow reuse) {
-    for (FieldSpec fieldSpec : 
_segmentMetadata.getSchema().getAllFieldSpecs()) {
+    Schema schema = _segmentMetadata.getSchema();
+    for (FieldSpec fieldSpec : schema.getAllFieldSpecs()) {
       String column = fieldSpec.getName();
-      ColumnIndexContainer indexContainer = _indexContainerMap.get(column);
-      reuse.putField(column, IndexSegmentUtils
-          .getValue(docId, fieldSpec, indexContainer.getForwardIndex(), 
indexContainer.getDictionary(),
-              
_segmentMetadata.getColumnMetadataFor(column).getMaxNumberOfMultiValues()));
+      if (!schema.isVirtualColumn(column)) {
+        ColumnIndexContainer indexContainer = _indexContainerMap.get(column);
+        reuse.putField(column, IndexSegmentUtils
+            .getValue(docId, fieldSpec, indexContainer.getForwardIndex(), 
indexContainer.getDictionary(),
+                
_segmentMetadata.getColumnMetadataFor(column).getMaxNumberOfMultiValues()));
+      }
     }
     return reuse;
   }
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImpl.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImpl.java
index 19eb51a..8591219 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImpl.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImpl.java
@@ -21,8 +21,12 @@ package org.apache.pinot.core.indexsegment.mutable;
 import com.google.common.base.Preconditions;
 import it.unimi.dsi.fastutil.ints.IntArrays;
 import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
@@ -95,6 +99,9 @@ public class MutableSegmentImpl implements MutableSegment {
   private volatile long _minTime = Long.MAX_VALUE;
   private volatile long _maxTime = Long.MIN_VALUE;
   private final int _numKeyColumns;
+  private final Collection<FieldSpec> _physicalColumnFieldSpecs;
+  private final Collection<FieldSpec> _physicalDimensionsFieldSpecs;
+  private final Collection<FieldSpec> _physicalMetricsFieldSpecs;
 
   // default message metadata
   private volatile long _lastIndexedTimeMs = Long.MIN_VALUE;
@@ -131,7 +138,27 @@ public class MutableSegmentImpl implements MutableSegment {
     _memoryManager = config.getMemoryManager();
     _statsHistory = config.getStatsHistory();
     _segmentPartitionConfig = config.getSegmentPartitionConfig();
-    _numKeyColumns = _schema.getDimensionNames().size() + 1;
+
+    List<FieldSpec> physicalColumnFieldSpecs = new 
ArrayList<>(_schema.getAllFieldSpecs());
+    Iterator<FieldSpec> iterator = physicalColumnFieldSpecs.iterator();
+    List<FieldSpec> physicalMetricsFieldSpecs = new 
ArrayList<>(_schema.getMetricNames().size());
+    List<FieldSpec> physicalDimensionsFieldSpecs = new 
ArrayList<>(_schema.getDimensionNames().size());
+    while (iterator.hasNext()) {
+      FieldSpec fieldSpec = iterator.next();
+      if (_schema.isVirtualColumn(fieldSpec.getName())) {
+        iterator.remove();
+      }
+      if (fieldSpec.getFieldType() == FieldSpec.FieldType.METRIC) {
+        physicalMetricsFieldSpecs.add(fieldSpec);
+      }
+      if (fieldSpec.getFieldType() == FieldSpec.FieldType.DIMENSION) {
+        physicalDimensionsFieldSpecs.add(fieldSpec);
+      }
+    }
+    _physicalColumnFieldSpecs = 
Collections.unmodifiableCollection(physicalColumnFieldSpecs);
+    _physicalDimensionsFieldSpecs = 
Collections.unmodifiableCollection(physicalDimensionsFieldSpecs);
+    _physicalMetricsFieldSpecs = 
Collections.unmodifiableCollection(physicalMetricsFieldSpecs);
+    _numKeyColumns = _physicalDimensionsFieldSpecs.size() + 1;  // Add 1 for 
time column
 
     _logger =
         LoggerFactory.getLogger(MutableSegmentImpl.class.getName() + "_" + 
_segmentName + "_" + config.getStreamName());
@@ -142,7 +169,7 @@ public class MutableSegmentImpl implements MutableSegment {
     int avgNumMultiValues = config.getAvgNumMultiValues();
 
     // Initialize for each column
-    for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
       String column = fieldSpec.getName();
       _maxNumValuesMap.put(column, 0);
 
@@ -197,7 +224,7 @@ public class MutableSegmentImpl implements MutableSegment {
 
     // Metric aggregation can be enabled only if config is specified, and all 
dimensions have dictionary,
     // and no metrics have dictionary. If not enabled, the map returned is 
null.
-    _recordIdMap = enableMetricsAggregationIfPossible(config, _schema, 
noDictionaryColumns);
+    _recordIdMap = enableMetricsAggregationIfPossible(config, 
noDictionaryColumns);
   }
 
   public SegmentPartitionConfig getSegmentPartitionConfig() {
@@ -250,7 +277,7 @@ public class MutableSegmentImpl implements MutableSegment {
 
   private Map<String, Object> updateDictionary(GenericRow row) {
     Map<String, Object> dictIdMap = new HashMap<>();
-    for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
       String column = fieldSpec.getName();
       Object value = row.getValue(column);
       MutableDictionary dictionary = _dictionaryMap.get(column);
@@ -294,7 +321,7 @@ public class MutableSegmentImpl implements MutableSegment {
 
   private void addForwardIndex(GenericRow row, int docId, Map<String, Object> 
dictIdMap) {
     // Store dictionary Id(s) for columns with dictionary
-    for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
       String column = fieldSpec.getName();
       Object value = row.getValue(column);
       if (fieldSpec.isSingleValueField()) {
@@ -336,7 +363,7 @@ public class MutableSegmentImpl implements MutableSegment {
     // Update inverted index at last
     // NOTE: inverted index have to be updated at last because once it gets 
updated, the latest record will become
     // queryable
-    for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
       String column = fieldSpec.getName();
       RealtimeInvertedIndexReader invertedIndex = 
_invertedIndexMap.get(column);
       if (invertedIndex != null) {
@@ -353,7 +380,7 @@ public class MutableSegmentImpl implements MutableSegment {
   }
 
   private boolean aggregateMetrics(GenericRow row, int docId) {
-    for (FieldSpec metricSpec : _schema.getMetricFieldSpecs()) {
+    for (FieldSpec metricSpec : _physicalMetricsFieldSpecs) {
       String column = metricSpec.getName();
       Object value = row.getValue(column);
       Preconditions.checkState(metricSpec.isSingleValueField(), "Multivalued 
metrics cannot be updated.");
@@ -401,6 +428,7 @@ public class MutableSegmentImpl implements MutableSegment {
 
   @Override
   public Set<String> getColumnNames() {
+    // Return all column names, virtual and physical.
     return _schema.getColumnNames();
   }
 
@@ -408,10 +436,8 @@ public class MutableSegmentImpl implements MutableSegment {
   public Set<String> getPhysicalColumnNames() {
     HashSet<String> physicalColumnNames = new HashSet<>();
 
-    for (String columnName : getColumnNames()) {
-      if (!_segmentMetadata.getSchema().isVirtualColumn(columnName)) {
-        physicalColumnNames.add(columnName);
-      }
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
+      physicalColumnNames.add(fieldSpec.getName());
     }
 
     return physicalColumnNames;
@@ -442,12 +468,17 @@ public class MutableSegmentImpl implements MutableSegment 
{
     return null;
   }
 
-  @Override
+  /**
+   * Returns a record that contains only physical columns
+   * @param docId document ID
+   * @param reuse a GenericRow object that will be re-used if provided. 
Otherwise, this method will allocate a new one
+   * @return Generic row with physical columns of the specified row.
+   */
   public GenericRow getRecord(int docId, GenericRow reuse) {
-    for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+    for (FieldSpec fieldSpec : _physicalColumnFieldSpecs) {
       String column = fieldSpec.getName();
-      reuse.putField(column, IndexSegmentUtils
-          .getValue(docId, fieldSpec, _indexReaderWriterMap.get(column), 
_dictionaryMap.get(column),
+      reuse.putField(column,
+          IndexSegmentUtils.getValue(docId, fieldSpec, 
_indexReaderWriterMap.get(column), _dictionaryMap.get(column),
               _maxNumValuesMap.getOrDefault(column, 0)));
     }
     return reuse;
@@ -573,8 +604,8 @@ public class MutableSegmentImpl implements MutableSegment {
     int[] dictIds = new int[_numKeyColumns]; // dimensions + time column.
 
     // FIXME: this for loop breaks for multi value dimensions. 
https://github.com/apache/incubator-pinot/issues/3867
-    for (String column : _schema.getDimensionNames()) {
-      dictIds[i++] = (Integer) dictIdMap.get(column);
+    for (FieldSpec fieldSpec : _physicalDimensionsFieldSpecs) {
+      dictIds[i++] = (Integer) dictIdMap.get(fieldSpec.getName());
     }
 
     String timeColumnName = _schema.getTimeColumnName();
@@ -591,18 +622,17 @@ public class MutableSegmentImpl implements MutableSegment 
{
    *   <li> All dimensions and time are dictionary encoded. This is because an 
integer array containing dictionary id's
    *        is used as key for dimensions to record Id map. </li>
    *   <li> None of the metrics are dictionary encoded. </li>
+   *   <li> All columns should be single-valued (see 
https://github.com/apache/incubator-pinot/issues/3867)</li>
    * </ul>
    *
    * TODO: Eliminate the requirement on dictionary encoding for dimension and 
metric columns.
    *
    * @param config Segment config.
-   * @param schema Schema for the table.
    * @param noDictionaryColumns Set of no dictionary columns.
    *
    * @return Map from dictionary id array to doc id, null if metrics 
aggregation cannot be enabled.
    */
-  private IdMap<FixedIntArray> 
enableMetricsAggregationIfPossible(RealtimeSegmentConfig config, Schema schema,
-      Set<String> noDictionaryColumns) {
+  private IdMap<FixedIntArray> 
enableMetricsAggregationIfPossible(RealtimeSegmentConfig config, Set<String> 
noDictionaryColumns) {
     _aggregateMetrics = config.aggregateMetrics();
     if (!_aggregateMetrics) {
       _logger.info("Metrics aggregation is disabled.");
@@ -611,15 +641,16 @@ public class MutableSegmentImpl implements MutableSegment 
{
 
     // All metric columns should have no-dictionary index.
     // All metric columns must be single value
-    for (String metric : schema.getMetricNames()) {
+    for (FieldSpec fieldSpec : _physicalMetricsFieldSpecs) {
+      String metric = fieldSpec.getName();
       if (!noDictionaryColumns.contains(metric)) {
         _logger
             .warn("Metrics aggregation cannot be turned ON in presence of 
dictionary encoded metrics, eg: {}", metric);
         _aggregateMetrics = false;
         break;
       }
-      // https://github.com/apache/incubator-pinot/issues/3867
-      if (!schema.getMetricSpec(metric).isSingleValueField()) {
+
+      if (!fieldSpec.isSingleValueField()) {
         _logger
             .warn("Metrics aggregation cannot be turned ON in presence of 
multi-value metric columns, eg: {}", metric);
         _aggregateMetrics = false;
@@ -629,15 +660,16 @@ public class MutableSegmentImpl implements MutableSegment 
{
 
     // All dimension columns should be dictionary encoded.
     // All dimension columns must be single value
-    for (String dimension : schema.getDimensionNames()) {
+    for (FieldSpec fieldSpec : _physicalDimensionsFieldSpecs) {
+      String dimension = fieldSpec.getName();
       if (noDictionaryColumns.contains(dimension)) {
         _logger
             .warn("Metrics aggregation cannot be turned ON in presence of 
no-dictionary dimensions, eg: {}", dimension);
         _aggregateMetrics = false;
         break;
       }
-      // https://github.com/apache/incubator-pinot/issues/3867
-      if (!schema.getDimensionSpec(dimension).isSingleValueField()) {
+
+      if (!fieldSpec.isSingleValueField()) {
         _logger.warn("Metrics aggregation cannot be turned ON in presence of 
multi-value dimension columns, eg: {}",
             dimension);
         _aggregateMetrics = false;
@@ -646,7 +678,7 @@ public class MutableSegmentImpl implements MutableSegment {
     }
 
     // Time column should be dictionary encoded.
-    String timeColumn = schema.getTimeColumnName();
+    String timeColumn = _schema.getTimeColumnName();
     if (noDictionaryColumns.contains(timeColumn)) {
       _logger
           .warn("Metrics aggregation cannot be turned ON in presence of 
no-dictionary time column, eg: {}", timeColumn);
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImplAggregateMetricsTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImplAggregateMetricsTest.java
index 167f8d2..7bd81d6 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImplAggregateMetricsTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/indexsegment/mutable/MutableSegmentImplAggregateMetricsTest.java
@@ -24,6 +24,7 @@ import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Map;
 import java.util.Random;
+import java.util.concurrent.TimeUnit;
 import org.apache.commons.lang.RandomStringUtils;
 import org.apache.pinot.common.data.FieldSpec;
 import org.apache.pinot.common.data.Schema;
@@ -40,6 +41,7 @@ public class MutableSegmentImplAggregateMetricsTest {
   private static final String DIMENSION_2 = "dim2";
   private static final String METRIC = "metric";
   private static final String METRIC_2 = "metric2";
+  private static final String TIME_COLUMN = "time";
   private static final String KEY_SEPARATOR = "\t\t";
   private static final int NUM_ROWS = 10001;
 
@@ -52,6 +54,7 @@ public class MutableSegmentImplAggregateMetricsTest {
             .addSingleValueDimension(DIMENSION_2, FieldSpec.DataType.STRING)
             .addMetric(METRIC, FieldSpec.DataType.LONG)
             .addMetric(METRIC_2, FieldSpec.DataType.FLOAT)
+            .addTime(TIME_COLUMN, TimeUnit.DAYS, FieldSpec.DataType.INT)
             .build();
     _mutableSegmentImpl = MutableSegmentImplTestUtils
         .createMutableSegmentImpl(schema, new 
HashSet<>(Arrays.asList(DIMENSION_1, METRIC, METRIC_2)),
@@ -73,9 +76,11 @@ public class MutableSegmentImplAggregateMetricsTest {
     Map<String, Float> expectedValuesFloat = new HashMap<>();
     StreamMessageMetadata defaultMetadata = new 
StreamMessageMetadata(System.currentTimeMillis());
     for (int i = 0; i < NUM_ROWS; i++) {
+      int daysSinceEpoch = random.nextInt(10);
       GenericRow row = new GenericRow();
       row.putField(DIMENSION_1, random.nextInt(10));
       row.putField(DIMENSION_2, 
stringValues[random.nextInt(stringValues.length)]);
+      row.putField(TIME_COLUMN, daysSinceEpoch);
       // Generate random int to prevent overflow
       long metricValue = random.nextInt();
       row.putField(METRIC, metricValue);
@@ -106,7 +111,8 @@ public class MutableSegmentImplAggregateMetricsTest {
   }
 
   private String buildKey(GenericRow row) {
-    return String.valueOf(row.getValue(DIMENSION_1)) + KEY_SEPARATOR + 
row.getValue(DIMENSION_2);
+    return String.valueOf(row.getValue(DIMENSION_1)) + KEY_SEPARATOR +
+        row.getValue(DIMENSION_2) + KEY_SEPARATOR + row.getValue(TIME_COLUMN);
   }
 
   @AfterClass


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to