java-topology/defects/pinot/patch/pinot-0002-partial-upsert-primary-key-list-contains.md
russell@unturf.com 25c2bafdee undf: assign 694-720; stamp patches; ruby-0003/elixir-0002/r-source-0002/victoria-metrics-0002
New UNDF assignments (693→720):
  elixir-0002 → UNDF-2026-000000698 (typespec used_type_pairs O(T²))
  r-source-0002 → UNDF-2026-000000711 (.walkClassGraph match dedup O(S²))
  ruby-0003 → UNDF-2026-000000712 (RubyGems dependent_gems O(N²×D))
  victoria-metrics-0002 → UNDF-2026-000000717 (MetricName tag-filter O(T×I))

Total: 720 UNDF assigned
2026-03-29 22:28:31 -04:00

4.8 KiB
Raw Blame History

UNDF: UNDF-2026-000000710

pinot-0002: PartialUpsertHandler + ColumnarMerger List.contains per column in hot upsert path

Classification

  • Severity: MEDIUM
  • CWE: CWE-407 (Algorithmic Complexity — Inefficient Algorithmic Complexity)
  • Component:
    • pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
    • pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/merger/PartialUpsertColumnarMerger.java
  • Methods: merge()

Defect

Both PartialUpsertHandler.merge() and PartialUpsertColumnarMerger.merge() iterate over all columns (C columns in the result holder / previous row) and for each column call _primaryKeyColumns.contains(column) and _comparisonColumns.contains(column). Both _primaryKeyColumns and _comparisonColumns are List<String>, so each .contains() is an O(P) or O(K) linear scan.

These merge() methods are called once per ingested row during streaming upsert. For a schema with C columns, P primary key columns, and K comparison columns, the per-row cost is O(C × (P + K)).

With C=100 columns, P=3 primary keys, and K=2 comparison columns this is 500 list scans per row. At 100k rows/sec throughput, that is 50 million redundant list scans per second.

Defective code — PartialUpsertHandler.java lines 48, 82

private final List<String> _primaryKeyColumns;   // line 48
private final List<String> _comparisonColumns;   // line 49

// merge() — called per row:
for (Map.Entry<String, Object> entry : resultHolder.entrySet()) {
    String column = entry.getKey();
    if (_primaryKeyColumns.contains(column)        // O(P) scan — per column per row
            || _comparisonColumns.contains(column)) { // O(K) scan — per column per row
        continue;
    }
    setMergedValue(newRow, column, entry.getValue());
}

Defective code — PartialUpsertColumnarMerger.java line 71

for (String column : previousRow.getColumnNames()) {
    if (_primaryKeyColumns.contains(column)        // O(P) per column per row
            || _comparisonColumns.contains(column)) { // O(K) per column per row
        continue;
    }
    ...
}

Fix

Build Set<String> lookups at construction time in both classes:

// PartialUpsertHandler constructor
private final Set<String> _primaryKeyColumnsSet;
private final Set<String> _comparisonColumnsSet;

// In constructor:
_primaryKeyColumnsSet = new HashSet<>(_primaryKeyColumns);
_comparisonColumnsSet = new HashSet<>(comparisonColumns);

// merge() — O(1) per column:
if (_primaryKeyColumnsSet.contains(column) || _comparisonColumnsSet.contains(column)) {
    continue;
}

Same pattern applies in PartialUpsertColumnarMerger (which already has _primaryKeyColumns and _comparisonColumns as List fields — add parallel Set fields).

Complexity

C columns P+K keys Before (per row) After (per row)
50 5 250 ops 50 ops
100 5 500 ops 100 ops
200 10 2000 ops 200 ops

Speedup: 5×10× on per-row merge path. At 100k rows/sec this directly reduces CPU consumed by partial upsert ingestion.

Patch

--- a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
+++ b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
@@ -1,3 +1,4 @@
+import java.util.HashSet;
+import java.util.Set;

   private final List<String> _primaryKeyColumns;
   private final List<String> _comparisonColumns;
+  private final Set<String> _primaryKeyColumnsSet;
+  private final Set<String> _comparisonColumnsSet;

   public PartialUpsertHandler(...) {
     _primaryKeyColumns = schema.getPrimaryKeyColumns();
     _comparisonColumns = comparisonColumns;
+    _primaryKeyColumnsSet = new HashSet<>(_primaryKeyColumns);
+    _comparisonColumnsSet = new HashSet<>(_comparisonColumns);
     ...
   }

   public void merge(...) {
     ...
     for (Map.Entry<String, Object> entry : resultHolder.entrySet()) {
       String column = entry.getKey();
-      if (_primaryKeyColumns.contains(column) || _comparisonColumns.contains(column)) {
+      if (_primaryKeyColumnsSet.contains(column) || _comparisonColumnsSet.contains(column)) {
         continue;
       }
       ...
     }
   }

--- a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/merger/PartialUpsertColumnarMerger.java
+++ b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/merger/PartialUpsertColumnarMerger.java
@@ apply same Set pattern to _primaryKeyColumns and _comparisonColumns fields
-      if (_primaryKeyColumns.contains(column) || _comparisonColumns.contains(column)) {
+      if (_primaryKeyColumnsSet.contains(column) || _comparisonColumnsSet.contains(column)) {