java-topology/defects/flink/patch/flink-0004-dynamicsinkutils-updatedcolumns-hashmap.md

2.9 KiB
Raw Blame History

UNDF: UNDF-2026-000000397

flink-0004: DynamicSinkUtils UPDATE column resolution O(C×U) → O(C+U)

Location

flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/DynamicSinkUtils.java

Severity

MEDIUM — triggered on every row-level UPDATE query plan compilation

Description

getUpdatedColumns() iterates over all schema columns and calls updatedColumnNames.contains(column.getName()) where updatedColumnNames is a List<String> from tableModify.getUpdateColumnList(). For a table with C columns and U updated columns, this is O(C×U).

projectColumnsForUpdate() iterates over updatedIndexes and calls both updatedColumnNames.contains(colName) (O(U)) and updatedColumnNames.indexOf(colName) (O(U)) per iteration — two linear scans per index, O(2×I×U) total.

Root Cause

// getUpdatedColumns — line 541-544
List<String> updatedColumnNames = tableModify.getUpdateColumnList(); // List<String>
for (Column column : resolvedSchema.getColumns()) {                  // O(C) loop
    if (updatedColumnNames.contains(column.getName())) {             // O(U) each → O(C×U)
        updatedColumns.add(column);
    }
}

// projectColumnsForUpdate — line 775-781
List<String> updatedColumnNames = tableModify.getUpdateColumnList(); // List<String>
for (int index : updatedIndexes) {                                   // O(I) loop
    String colName = resolvedSchema.getColumnNames().get(index);
    if (updatedColumnNames.contains(colName)) {                      // O(U) → O(I×U)
        int i = updatedColumnNames.indexOf(colName);                 // O(U) again!
        ...
    }
}

Fix

Pre-build a Set<String> for O(1) contains, and a Map<String,Integer> for O(1) indexOf.

// getUpdatedColumns fix
Set<String> updatedColumnSet = new HashSet<>(tableModify.getUpdateColumnList());
for (Column column : resolvedSchema.getColumns()) {
    if (updatedColumnSet.contains(column.getName())) {   // O(1)
        updatedColumns.add(column);
    }
}

// projectColumnsForUpdate fix
List<String> updatedColumnNames = tableModify.getUpdateColumnList();
Map<String, Integer> updatedColumnIndex = new HashMap<>();
for (int i = 0; i < updatedColumnNames.size(); i++) {
    updatedColumnIndex.put(updatedColumnNames.get(i), i);
}
for (int index : updatedIndexes) {
    String colName = resolvedSchema.getColumnNames().get(index);
    Integer i = updatedColumnIndex.get(colName);    // O(1) replaces contains+indexOf
    if (i != null) {
        RexNode rexNode = oldRexNodes.get(originColsCount + i);
        ...
    }
}

Complexity

Before After
getUpdatedColumns O(C×U) O(C+U)
projectColumnsForUpdate O(2×I×U) O(I+U)

Where C=total columns, U=updated columns, I=updated index count. At C=500 columns, U=50 updates: 25,000 ops → 550 ops (45x improvement).

Duplicate locations

Both methods are in the same file. No other copies found.