From 6be413b510af45990aff9cdf9645d8410e218bd8 Mon Sep 17 00:00:00 2001 From: sonika-shah <58761340+sonika-shah@users.noreply.github.com> Date: Tue, 8 Sep 2026 02:23:31 +0530 Subject: [PATCH 1/4] Fixes #26824: filter column-grid metadata status on the aggregate row, not per document The "Has / Missing Metadata" filter in Column Bulk Operations did not hide rows of other statuses, and page counts drifted per page. metadataStatus (MISSING/INCOMPLETE/COMPLETE) was pushed down as a per-document search query, while INCONSISTENT/hasConflicts/hasMissingMetadata were applied after the aggregator had already paginated and computed totals. But a row's status is an aggregate over all of a column's occurrences (INCONSISTENT when they disagree), so a per-document filter re-grouped into rows whose status differed from the request, and the post-pagination filter shrank the page while the total stayed unfiltered. Move every row-level filter (metadataStatus, hasConflicts, hasMissingMetadata) onto the fully-grouped items, before pagination, with totals derived from the filtered set. Shared pure helpers (hasRowLevelFilter, matchesRowFilters, paginateFilteredItems) live on the ColumnAggregator interface; both aggregators route through a materialized name-enumeration path when such a filter is active, and the untouched composite/pattern paths still serve unfiltered browsing. Removing the per-document status query also removes the wildcard/exists query on flat-object columns.description/columns.tags that crashed ES/OS with search_phase_execution_exception and had ColumnGridResourceIT disabled; the IT is re-enabled and its status assertions strengthened to check every returned row carries the requested status, plus a status+pagination consistency test. --- .../it/tests/ColumnGridResourceIT.java | 101 +++++++++++-- .../service/jdbi3/ColumnRepository.java | 38 +---- .../resources/columns/ColumnResource.java | 8 +- .../service/search/ColumnAggregator.java | 98 +++++++++++++ .../ElasticSearchColumnAggregator.java | 119 ++++++++-------- .../OpenSearchColumnAggregator.java | 111 +++++++-------- .../service/search/ColumnAggregatorTest.java | 134 ++++++++++++++++++ 7 files changed, 445 insertions(+), 164 deletions(-) diff --git a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java index 8d597cdbfb50..c4790fc4aee1 100644 --- a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java +++ b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java @@ -13,7 +13,6 @@ import java.time.Duration; import java.util.List; import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.parallel.Execution; @@ -53,17 +52,10 @@ import org.openmetadata.sdk.fluent.Tables; import org.openmetadata.sdk.network.HttpMethod; -// TEMPORARILY DISABLED — the metadataStatus aggregation on this endpoint reproducibly fails -// with [search_phase_execution_exception] all shards failed on both postgres+ES+redis (single -// failure on test_getColumnGrid_withMetadataStatusIncomplete) AND postgres+OpenSearch (the same -// query crashes the OS container, then 15 follow-up tests in the class fail with Connection -// refused). Same behavior on PR #28100 with and without the cache changes, so it is a -// pre-existing aggregator bug, not a cache regression. The ES Java client swallows the -// underlying `caused_by`, so root-causing the actual ES-side error requires response-body -// logging that is not wired up yet. Re-enable once the underlying aggregator/index-mapping -// issue is fixed in a follow-up. See PR #28100 history and CI run 25940411417 for context. -@Disabled( - "ColumnGrid metadataStatus aggregation crashes ES/OS — pre-existing flake, follow-up needed") +// Re-enabled with #26824: the metadataStatus crash came from the per-document filter query +// (wildcard/exists on flat-object columns.description/columns.tags), which ES 7.17 and OpenSearch +// rejected with `search_phase_execution_exception ... all shards failed`. That push-down is gone — +// status is now filtered on the aggregate grouped item — so the crashing query no longer runs. @Execution(ExecutionMode.CONCURRENT) @ExtendWith(TestNamespaceExtension.class) public class ColumnGridResourceIT { @@ -314,6 +306,11 @@ void test_getColumnGrid_withMetadataStatusMissing(TestNamespace ns) throws Excep assertNotNull(response); assertNotNull(response.getColumns()); + assertAllRowsHaveStatus(response, MetadataStatus.MISSING); + assertEquals( + response.getColumns().size(), + response.getTotalUniqueColumns(), + "totalUniqueColumns must reflect the filtered set, not the unfiltered total"); } @Test @@ -328,6 +325,9 @@ void test_getColumnGrid_withMetadataStatusComplete(TestNamespace ns) throws Exce assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the COMPLETE column should be returned"); + // The reported bug (#26824): COMPLETE must not surface MISSING/INCOMPLETE/INCONSISTENT rows. + assertAllRowsHaveStatus(response, MetadataStatus.COMPLETE); } @Test @@ -342,6 +342,8 @@ void test_getColumnGrid_withMetadataStatusIncomplete(TestNamespace ns) throws Ex assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the INCOMPLETE column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.INCOMPLETE); } @Test @@ -422,6 +424,80 @@ void test_getColumnGrid_withMetadataStatusInconsistent(TestNamespace ns) throws assertNotNull(response); assertNotNull(response.getColumns()); + assertFalse(response.getColumns().isEmpty(), "the INCONSISTENT column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.INCONSISTENT); + assertTrue( + response.getColumns().stream().allMatch(ColumnGridItem::getHasVariations), + "INCONSISTENT rows have metadata variations across occurrences"); + } + + @Test + void test_getColumnGrid_metadataStatusPaginationCountsAreConsistent(TestNamespace ns) + throws Exception { + OpenMetadataClient client = SdkClients.adminClient(); + DatabaseService service = DatabaseServiceTestFactory.createPostgres(ns); + DatabaseSchema schema = DatabaseSchemaTestFactory.createSimple(ns, service); + + // Three COMPLETE columns and one MISSING column in the same service. + for (int i = 0; i < 3; i++) { + Column complete = + Columns.build(ns.prefix("paged_complete_" + i)) + .withType(ColumnDataType.BIGINT) + .withDescription("has description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) + .create(); + Tables.create() + .name(ns.prefix("paged_table_" + i)) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(complete)) + .execute(); + } + Column missing = + Columns.build(ns.prefix("paged_missing")).withType(ColumnDataType.BIGINT).create(); + Tables.create() + .name(ns.prefix("paged_table_missing")) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(missing)) + .execute(); + + waitForSearchIndexRefresh(ns); + + ColumnGridResponse page1 = + getColumnGrid( + client, + "size=2&entityTypes=table&metadataStatus=COMPLETE&serviceName=" + service.getName()); + + // totalUniqueColumns must count only the 3 COMPLETE columns (not 4), and the page must respect + // the requested size — the pagination half of #26824. + assertEquals(3, page1.getTotalUniqueColumns()); + assertEquals(2, page1.getColumns().size()); + assertAllRowsHaveStatus(page1, MetadataStatus.COMPLETE); + assertNotNull(page1.getCursor(), "a second page of COMPLETE columns remains"); + + ColumnGridResponse page2 = + getColumnGrid( + client, + "size=2&entityTypes=table&metadataStatus=COMPLETE&serviceName=" + + service.getName() + + "&cursor=" + + URLEncoder.encode(page1.getCursor(), StandardCharsets.UTF_8)); + + assertEquals(3, page2.getTotalUniqueColumns()); + assertEquals(1, page2.getColumns().size(), "the last page holds the remaining COMPLETE column"); + assertAllRowsHaveStatus(page2, MetadataStatus.COMPLETE); + } + + /** + * Every returned row must carry the requested aggregate status — the core guarantee of #26824 + * (before the fix a COMPLETE/INCOMPLETE filter leaked rows of other statuses). + */ + private void assertAllRowsHaveStatus(ColumnGridResponse response, MetadataStatus expected) { + for (ColumnGridItem item : response.getColumns()) { + assertEquals( + expected, + item.getMetadataStatus(), + "column '" + item.getColumnName() + "' should have status " + expected); + } } @Test @@ -1424,6 +1500,7 @@ private DatabaseService createTableWithFullMetadata(TestNamespace ns) { Columns.build("full_metadata_id") .withType(ColumnDataType.BIGINT) .withDescription("Primary key with description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) .create(); Tables.create() diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java b/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java index 88373c992ac0..518a3930eb6c 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/jdbi3/ColumnRepository.java @@ -40,7 +40,6 @@ import org.openmetadata.schema.EntityInterface; import org.openmetadata.schema.api.data.BulkColumnUpdatePreview; import org.openmetadata.schema.api.data.BulkColumnUpdateRequest; -import org.openmetadata.schema.api.data.ColumnGridItem; import org.openmetadata.schema.api.data.ColumnGridResponse; import org.openmetadata.schema.api.data.ColumnMetadata; import org.openmetadata.schema.api.data.ColumnOccurrence; @@ -97,39 +96,10 @@ public ColumnRepository(Authorizer authorizer, SearchClient searchClient) { public ColumnGridResponse getColumnGridPaginated( SecurityContext securityContext, ColumnAggregator.ColumnAggregationRequest request) throws IOException { - ColumnGridResponse response = columnAggregator.aggregateColumns(request); - - if (Boolean.TRUE.equals(request.getHasConflicts())) { - response.setColumns( - response.getColumns().stream() - .filter(ColumnGridItem::getHasVariations) - .collect(Collectors.toList())); - } - - if (Boolean.TRUE.equals(request.getHasMissingMetadata())) { - response.setColumns( - response.getColumns().stream() - .filter(this::hasMissingMetadata) - .collect(Collectors.toList())); - } - - // Filter by INCONSISTENT status (requires post-aggregation filtering) - if ("INCONSISTENT".equalsIgnoreCase(request.getMetadataStatus())) { - response.setColumns( - response.getColumns().stream() - .filter(ColumnGridItem::getHasVariations) - .collect(Collectors.toList())); - } - - return response; - } - - private boolean hasMissingMetadata(ColumnGridItem item) { - return item.getGroups().stream() - .anyMatch( - group -> - (group.getDescription() == null || group.getDescription().isEmpty()) - || (group.getTags() == null || group.getTags().isEmpty())); + // Row-level filters (metadataStatus / hasConflicts / hasMissingMetadata) are applied inside the + // aggregator over the fully-grouped columns, before pagination, so page counts and per-page + // size stay correct (#26824). Nothing to post-process here. + return columnAggregator.aggregateColumns(request); } public Column getColumnByFQN( diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java b/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java index cd40605a90b3..65103fecd31c 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/resources/columns/ColumnResource.java @@ -319,13 +319,15 @@ public Response getColumnGrid( boolean hasMissingMetadata, @Parameter( description = - "Filter by metadata status: MISSING (no description AND no tags), " + "Filter by aggregate metadata status of a column across all its occurrences: " + + "MISSING (no description AND no tags), " + "INCOMPLETE (has description OR tags, but not both), " - + "COMPLETE (has both description AND tags)", + + "COMPLETE (has both description AND tags), " + + "INCONSISTENT (occurrences disagree on description/tags)", schema = @Schema( type = "string", - allowableValues = {"MISSING", "INCOMPLETE", "COMPLETE"})) + allowableValues = {"MISSING", "INCOMPLETE", "COMPLETE", "INCONSISTENT"})) @QueryParam("metadataStatus") String metadataStatus, @Parameter( diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java index f673b20c33e2..5106dafd549d 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java @@ -16,9 +16,12 @@ import com.fasterxml.jackson.core.type.TypeReference; import java.io.IOException; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.Base64; +import java.util.Comparator; import java.util.List; import java.util.Map; +import org.openmetadata.schema.api.data.ColumnGridItem; import org.openmetadata.schema.api.data.ColumnGridResponse; import org.openmetadata.schema.utils.JsonUtils; import org.slf4j.Logger; @@ -31,6 +34,15 @@ public interface ColumnAggregator { /** Max column names to retrieve in the names-only query during pattern search. */ int MAX_PATTERN_SEARCH_NAMES = 10000; + /** Terms include array matching every column name (row-level-filter path has no name pattern). */ + String MATCH_ALL_NAMES_REGEX = ".*"; + + /** + * Batch size for fetching occurrence data by column name in the row-level-filter path. Bounds the + * terms-include array and top_hits fan-out per query when many names must be materialized. + */ + int ROW_FILTER_DATA_BATCH = 1000; + /** * Number of sample docs pulled per column-name bucket to populate occurrences. Caps * {@code ColumnGridItem.totalOccurrences}; columns appearing in more entities than this @@ -115,6 +127,92 @@ static int toIntSaturating(long value) { return (int) value; } + /** + * True if a request carries a row-level filter — one that acts on the aggregate status of a + * grouped column (metadataStatus, hasConflicts, hasMissingMetadata) rather than on individual + * documents. These cannot be pushed into the search query (the aggregate status is only known + * after grouping occurrences), so they are applied via {@link #paginateFilteredItems}. + */ + static boolean hasRowLevelFilter(ColumnAggregationRequest request) { + return !nullOrEmptyStr(request.getMetadataStatus()) + || Boolean.TRUE.equals(request.getHasConflicts()) + || Boolean.TRUE.equals(request.getHasMissingMetadata()); + } + + /** True if the grouped column satisfies every active row-level filter on the request. */ + static boolean matchesRowFilters(ColumnGridItem item, ColumnAggregationRequest request) { + if (Boolean.TRUE.equals(request.getHasConflicts()) + && !Boolean.TRUE.equals(item.getHasVariations())) { + return false; + } + if (Boolean.TRUE.equals(request.getHasMissingMetadata()) && !itemHasMissingMetadata(item)) { + return false; + } + String status = request.getMetadataStatus(); + if (!nullOrEmptyStr(status)) { + String itemStatus = item.getMetadataStatus() != null ? item.getMetadataStatus().value() : ""; + return status.trim().equalsIgnoreCase(itemStatus); + } + return true; + } + + /** A column has missing metadata if any of its groups lacks a description or tags. */ + static boolean itemHasMissingMetadata(ColumnGridItem item) { + if (item.getGroups() == null) { + return true; + } + return item.getGroups().stream() + .anyMatch( + group -> + (group.getDescription() == null || group.getDescription().isEmpty()) + || (group.getTags() == null || group.getTags().isEmpty())); + } + + /** + * Apply row-level filters to the fully-grouped column list and paginate the result in memory. + * Filtering the aggregate items (not documents) is what makes a "Complete"/"Incomplete"/… filter + * return only rows whose displayed status matches, and computing totals from the filtered set is + * what keeps the page count and per-page size correct (issue #26824). Ordering is by column name + * (case-insensitive) so the offset cursor is stable across pages. + */ + static ColumnGridResponse paginateFilteredItems( + List allItems, ColumnAggregationRequest request) { + List filtered = + allItems.stream() + .filter(item -> matchesRowFilters(item, request)) + .sorted( + Comparator.comparing( + ColumnGridItem::getColumnName, + Comparator.nullsLast(String.CASE_INSENSITIVE_ORDER))) + .toList(); + + int totalUniqueColumns = filtered.size(); + int totalOccurrences = + filtered.stream() + .mapToInt(item -> item.getTotalOccurrences() != null ? item.getTotalOccurrences() : 0) + .sum(); + + int offset = decodeSearchOffset(request.getCursor()); + int pageSize = request.getSize(); + int fromIndex = Math.min(offset, totalUniqueColumns); + int toIndex = Math.min(offset + pageSize, totalUniqueColumns); + + List page = new ArrayList<>(filtered.subList(fromIndex, toIndex)); + boolean hasMore = toIndex < totalUniqueColumns; + String cursor = hasMore ? encodeSearchOffset(toIndex) : null; + + ColumnGridResponse response = new ColumnGridResponse(); + response.setColumns(page); + response.setTotalUniqueColumns(totalUniqueColumns); + response.setTotalOccurrences(totalOccurrences); + response.setCursor(cursor); + return response; + } + + private static boolean nullOrEmptyStr(String s) { + return s == null || s.isBlank(); + } + /** Phase 1 result: matching column names and the total doc_count summed across buckets. */ record NamesWithCount(List names, long totalDocCount) {} diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java index f91640f948be..b1b2d63fee42 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java @@ -115,9 +115,22 @@ public ColumnGridResponse aggregateColumns(ColumnAggregationRequest request) thr .removeIf(e -> !e.getKey().toLowerCase(Locale.ROOT).contains(pattern)); } + // Row-level filter (metadataStatus / hasConflicts / hasMissingMetadata) acts on the + // aggregate status, so group everything then filter + paginate the items. + if (ColumnAggregator.hasRowLevelFilter(request)) { + List gridItems = ColumnMetadataGrouper.groupColumns(taggedColumns); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + return aggregateColumnsWithKnownNames(request, taggedColumns); } + // Row-level filter (no tag filter): materialize all candidate columns, then filter the + // aggregate items and paginate in memory so counts and per-page size stay correct (#26824). + if (ColumnAggregator.hasRowLevelFilter(request)) { + return aggregateColumnsWithRowFilters(request, entityTypes); + } + // Pattern-only path (no tag filter): use terms agg with include regex if (!nullOrEmpty(request.getColumnNamePattern())) { return aggregateColumnsWithPattern(request, entityTypes); @@ -266,6 +279,54 @@ private ColumnGridResponse aggregateColumnsWithPattern( return buildResponse(gridItems, cursor, hasMore, totalUniqueColumns, totalOccurrences); } + /** + * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The + * filter acts on the aggregate status of a grouped column, which is only known after grouping all + * of its occurrences — so we enumerate every candidate name (respecting any column-name pattern + * and scope filters), fetch their occurrences, group them, then filter + paginate the items in + * memory. This keeps the page count and per-page size consistent with the filtered result set. + */ + private ColumnGridResponse aggregateColumnsWithRowFilters( + ColumnAggregationRequest request, List entityTypes) throws IOException { + + Map> fieldPathToEntityTypes = groupByFieldPath(entityTypes); + String regex = + !nullOrEmpty(request.getColumnNamePattern()) + ? ColumnAggregator.toCaseInsensitiveRegex(request.getColumnNamePattern()) + : MATCH_ALL_NAMES_REGEX; + + Map> allColumnsByName = new HashMap<>(); + + for (Map.Entry> entry : fieldPathToEntityTypes.entrySet()) { + String columnNameKeyword = entry.getKey(); + List indexes = resolveIndexNames(entry.getValue()); + String columnFieldPath = INDEX_CONFIGS.get(entry.getValue().getFirst()).columnFieldPath(); + Query query = buildFilters(request, columnNameKeyword, null); + + try { + List names = executeNamesQuery(query, indexes, columnNameKeyword, regex).names(); + for (int i = 0; i < names.size(); i += ROW_FILTER_DATA_BATCH) { + List batch = names.subList(i, Math.min(i + ROW_FILTER_DATA_BATCH, names.size())); + Map> columnsByName = + executePageDataQuery(query, indexes, columnNameKeyword, columnFieldPath, batch); + for (Map.Entry> colEntry : columnsByName.entrySet()) { + allColumnsByName + .computeIfAbsent(colEntry.getKey(), k -> new ArrayList<>()) + .addAll(colEntry.getValue()); + } + } + } catch (ElasticsearchException e) { + if (!isIndexNotFoundException(e)) { + logShardFailureDetails(e, indexes, query); + throw e; + } + } + } + + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -461,7 +522,6 @@ private Query buildTagFilterQuery(ColumnAggregationRequest request, String colum addSchemaFilter(boolBuilder, request); addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request, columnNameKeyword); - addMetadataStatusFilter(boolBuilder, request, columnFieldPath); String tagFQNField = columnNameKeyword.replace(".name.keyword", ".tags.tagFQN"); List allTags = new ArrayList<>(); @@ -562,7 +622,6 @@ private Query buildFilters( addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request, columnNameKeyword); addTagFilters(boolBuilder, request, columnNameKeyword, columnNamesFromTagFilter); - addMetadataStatusFilter(boolBuilder, request, columnFieldPath); return Query.of(q -> q.bool(boolBuilder.build())); } @@ -680,62 +739,6 @@ private void addTagFilters( } } - private void addMetadataStatusFilter( - BoolQuery.Builder boolBuilder, ColumnAggregationRequest request, String columnFieldPath) { - if (!nullOrEmpty(request.getMetadataStatus())) { - Query metadataStatusQuery = - buildMetadataStatusFilter(request.getMetadataStatus(), columnFieldPath); - if (metadataStatusQuery != null) { - boolBuilder.filter(metadataStatusQuery); - } - } - } - - private Query buildMetadataStatusFilter(String status, String columnFieldPath) { - String descField = columnFieldPath + ".description"; - String tagsField = columnFieldPath + ".tags"; - - Query hasDesc = hasNonEmptyField(descField); - Query hasTags = existsQuery(tagsField); - Query noDesc = hasEmptyOrMissingField(descField); - Query noTags = notExistsQuery(tagsField); - - return switch (status.toUpperCase()) { - case "MISSING" -> Query.of(q -> q.bool(b -> b.must(noDesc).must(noTags))); - case "INCOMPLETE" -> Query.of( - q -> - q.bool( - b -> - b.should(Query.of(qs -> qs.bool(bs -> bs.must(hasDesc).must(noTags)))) - .should(Query.of(qs -> qs.bool(bs -> bs.must(noDesc).must(hasTags)))) - .minimumShouldMatch("1"))); - case "COMPLETE" -> Query.of(q -> q.bool(b -> b.must(hasDesc).must(hasTags))); - default -> null; - }; - } - - private Query existsQuery(String field) { - return Query.of(q -> q.exists(e -> e.field(field))); - } - - private Query notExistsQuery(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(existsQuery(field)))); - } - - // `wildcard(field, "?*")` matches any doc whose indexed terms include at least one token of - // at least one character — the analyzer-friendly equivalent of "field has non-empty value". - // We can't use `term(field, "")` against analyzed text fields like `columns.description`: the - // field's analyzer produces no tokens for the empty string and ES 7.17 rejects the term query - // with `search_phase_execution_exception ... all shards failed`. Caught by - // ColumnGridResourceIT#test_getColumnGrid_withMetadataStatusIncomplete. - private Query hasNonEmptyField(String field) { - return Query.of(q -> q.wildcard(w -> w.field(field).value("?*"))); - } - - private Query hasEmptyOrMissingField(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(hasNonEmptyField(field)))); - } - /** Phase 1: Get all matching column names using terms agg with include regex (no top_hits). */ private ColumnAggregator.NamesWithCount executeNamesQuery( Query query, List indexes, String columnNameKeyword, String regex) diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java index faac76576f7d..1c348bf1361d 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java @@ -98,9 +98,22 @@ public ColumnGridResponse aggregateColumns(ColumnAggregationRequest request) thr .removeIf(e -> !e.getKey().toLowerCase(Locale.ROOT).contains(pattern)); } + // Row-level filter (metadataStatus / hasConflicts / hasMissingMetadata) acts on the + // aggregate status, so group everything then filter + paginate the items. + if (ColumnAggregator.hasRowLevelFilter(request)) { + List gridItems = ColumnMetadataGrouper.groupColumns(taggedColumns); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } + return aggregateColumnsWithKnownNames(request, taggedColumns); } + // Row-level filter (no tag filter): materialize all candidate columns, then filter the + // aggregate items and paginate in memory so counts and per-page size stay correct (#26824). + if (ColumnAggregator.hasRowLevelFilter(request)) { + return aggregateColumnsWithRowFilters(request); + } + // Pattern-only path (no tag filter): use terms agg with include regex if (!nullOrEmpty(request.getColumnNamePattern())) { return aggregateColumnsWithPattern(request); @@ -187,6 +200,47 @@ private ColumnGridResponse aggregateColumnsWithPattern(ColumnAggregationRequest } } + /** + * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The + * filter acts on the aggregate status of a grouped column, which is only known after grouping all + * of its occurrences — so we enumerate every candidate name (respecting any column-name pattern + * and scope filters), fetch their occurrences, group them, then filter + paginate the items in + * memory. This keeps the page count and per-page size consistent with the filtered result set. + */ + private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationRequest request) + throws IOException { + + Query query = buildFilters(request, null); + String regex = + !nullOrEmpty(request.getColumnNamePattern()) + ? ColumnAggregator.toCaseInsensitiveRegex(request.getColumnNamePattern()) + : ColumnAggregator.MATCH_ALL_NAMES_REGEX; + + try { + List names = executeNamesQuery(query, regex).names(); + Map> allColumnsByName = new HashMap<>(); + for (int i = 0; i < names.size(); i += ColumnAggregator.ROW_FILTER_DATA_BATCH) { + List batch = + names.subList(i, Math.min(i + ColumnAggregator.ROW_FILTER_DATA_BATCH, names.size())); + Map> columnsByName = executePageDataQuery(query, batch); + for (Map.Entry> colEntry : columnsByName.entrySet()) { + allColumnsByName + .computeIfAbsent(colEntry.getKey(), k -> new ArrayList<>()) + .addAll(colEntry.getValue()); + } + } + + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); + return ColumnAggregator.paginateFilteredItems(gridItems, request); + } catch (OpenSearchException e) { + if (isIndexNotFoundException(e)) { + LOG.warn("Search index not found, returning empty results"); + return buildResponse(new ArrayList<>(), null, false, 0, 0); + } + throw e; + } + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -376,7 +430,6 @@ private Query buildTagFilterQuery(ColumnAggregationRequest request) { addSchemaFilter(boolBuilder, request); addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request); - addMetadataStatusFilter(boolBuilder, request); List allTags = new ArrayList<>(); if (!nullOrEmpty(request.getTags())) { @@ -427,7 +480,6 @@ private Query buildFilters( addDomainFilter(boolBuilder, request); addColumnNamePatternFilter(boolBuilder, request); addTagFilters(boolBuilder, request, columnNamesFromTagFilter); - addMetadataStatusFilter(boolBuilder, request); return Query.of(q -> q.bool(boolBuilder.build())); } @@ -550,61 +602,6 @@ private void addTagFilters( } } - private void addMetadataStatusFilter( - BoolQuery.Builder boolBuilder, ColumnAggregationRequest request) { - if (!nullOrEmpty(request.getMetadataStatus())) { - Query metadataStatusQuery = buildMetadataStatusFilter(request.getMetadataStatus()); - if (metadataStatusQuery != null) { - boolBuilder.filter(metadataStatusQuery); - } - } - } - - private Query buildMetadataStatusFilter(String status) { - String descField = "columns.description"; - String tagsField = "columns.tags"; - - Query hasDesc = hasNonEmptyField(descField); - Query hasTags = existsQuery(tagsField); - Query noDesc = hasEmptyOrMissingField(descField); - Query noTags = notExistsQuery(tagsField); - - return switch (status.toUpperCase()) { - case "MISSING" -> Query.of(q -> q.bool(b -> b.must(noDesc).must(noTags))); - case "INCOMPLETE" -> Query.of( - q -> - q.bool( - b -> - b.should(Query.of(qs -> qs.bool(bs -> bs.must(hasDesc).must(noTags)))) - .should(Query.of(qs -> qs.bool(bs -> bs.must(noDesc).must(hasTags)))) - .minimumShouldMatch("1"))); - case "COMPLETE" -> Query.of(q -> q.bool(b -> b.must(hasDesc).must(hasTags))); - default -> null; - }; - } - - private Query existsQuery(String field) { - return Query.of(q -> q.exists(e -> e.field(field))); - } - - private Query notExistsQuery(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(existsQuery(field)))); - } - - // `wildcard(field, "?*")` matches any doc whose indexed terms include at least one token of - // at least one character — the analyzer-friendly equivalent of "field has non-empty value". - // We can't use `term(field, "")` against analyzed text fields like `columns.description`: the - // field's analyzer produces no tokens for the empty string and OS rejects the term query with - // `search_phase_execution_exception ... all shards failed`. Caught by - // ColumnGridResourceIT#test_getColumnGrid_withMetadataStatusIncomplete. - private Query hasNonEmptyField(String field) { - return Query.of(q -> q.wildcard(w -> w.field(field).value("?*"))); - } - - private Query hasEmptyOrMissingField(String field) { - return Query.of(q -> q.bool(b -> b.mustNot(hasNonEmptyField(field)))); - } - /** Phase 1: Get all matching column names using terms agg with include regex (no top_hits). */ private ColumnAggregator.NamesWithCount executeNamesQuery(Query query, String regex) throws IOException { diff --git a/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java b/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java index b0f8259f010e..13c79adc3239 100644 --- a/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java +++ b/openmetadata-service/src/test/java/org/openmetadata/service/search/ColumnAggregatorTest.java @@ -15,10 +15,20 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; +import java.util.ArrayList; +import java.util.List; import java.util.regex.Pattern; import org.junit.jupiter.api.Test; +import org.openmetadata.schema.api.data.ColumnGridItem; +import org.openmetadata.schema.api.data.ColumnGridResponse; +import org.openmetadata.schema.api.data.ColumnMetadataGroup; +import org.openmetadata.schema.api.data.MetadataStatus; +import org.openmetadata.schema.type.TagLabel; +import org.openmetadata.service.search.ColumnAggregator.ColumnAggregationRequest; class ColumnAggregatorTest { @@ -109,4 +119,128 @@ void toCaseInsensitiveRegex_specialCharsAreEscaped() { // Plus and star should be literal, not regex quantifiers assertFalse(pattern.matcher("abbbbc").matches()); } + + // ---- Row-level (post-grouping) filter + pagination helpers (issue #26824) ---- + + private static ColumnGridItem item(String name, MetadataStatus status, boolean hasVariations) { + ColumnGridItem gi = new ColumnGridItem(); + gi.setColumnName(name); + gi.setMetadataStatus(status); + gi.setHasVariations(hasVariations); + gi.setTotalOccurrences(1); + ColumnMetadataGroup group = new ColumnMetadataGroup(); + if (status == MetadataStatus.COMPLETE || status == MetadataStatus.INCONSISTENT) { + group.setDescription("desc"); + group.setTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))); + } else if (status == MetadataStatus.INCOMPLETE) { + group.setDescription("desc"); + } + gi.setGroups(new ArrayList<>(List.of(group))); + return gi; + } + + private static ColumnAggregationRequest request(String metadataStatus) { + ColumnAggregationRequest r = new ColumnAggregationRequest(); + r.setMetadataStatus(metadataStatus); + return r; + } + + @Test + void matchesRowFilters_metadataStatusMatchesAggregateStatus() { + ColumnGridItem complete = item("a", MetadataStatus.COMPLETE, false); + ColumnGridItem incomplete = item("b", MetadataStatus.INCOMPLETE, false); + + assertTrue(ColumnAggregator.matchesRowFilters(complete, request("COMPLETE"))); + assertFalse(ColumnAggregator.matchesRowFilters(incomplete, request("COMPLETE"))); + assertTrue(ColumnAggregator.matchesRowFilters(incomplete, request("INCOMPLETE"))); + // Case-insensitive + assertTrue(ColumnAggregator.matchesRowFilters(complete, request("complete"))); + } + + @Test + void matchesRowFilters_inconsistentIsAFirstClassStatus() { + ColumnGridItem inconsistent = item("a", MetadataStatus.INCONSISTENT, true); + ColumnGridItem complete = item("b", MetadataStatus.COMPLETE, false); + + // The reported bug: "COMPLETE" must NOT return inconsistent rows. + assertFalse(ColumnAggregator.matchesRowFilters(inconsistent, request("COMPLETE"))); + assertTrue(ColumnAggregator.matchesRowFilters(inconsistent, request("INCONSISTENT"))); + assertFalse(ColumnAggregator.matchesRowFilters(complete, request("INCONSISTENT"))); + } + + @Test + void matchesRowFilters_nullOrBlankStatusMatchesEverything() { + ColumnGridItem complete = item("a", MetadataStatus.COMPLETE, false); + assertTrue(ColumnAggregator.matchesRowFilters(complete, request(null))); + assertTrue(ColumnAggregator.matchesRowFilters(complete, request(" "))); + } + + @Test + void matchesRowFilters_hasConflictsRequiresVariations() { + ColumnAggregationRequest r = new ColumnAggregationRequest(); + r.setHasConflicts(true); + assertTrue(ColumnAggregator.matchesRowFilters(item("a", MetadataStatus.INCONSISTENT, true), r)); + assertFalse(ColumnAggregator.matchesRowFilters(item("b", MetadataStatus.COMPLETE, false), r)); + } + + @Test + void paginateFilteredItems_filtersAndComputesTotalsFromFilteredSet() { + List all = + new ArrayList<>( + List.of( + item("complete_1", MetadataStatus.COMPLETE, false), + item("inconsistent_1", MetadataStatus.INCONSISTENT, true), + item("complete_2", MetadataStatus.COMPLETE, false), + item("missing_1", MetadataStatus.MISSING, false))); + + ColumnAggregationRequest r = request("COMPLETE"); + r.setSize(10); + + ColumnGridResponse resp = ColumnAggregator.paginateFilteredItems(all, r); + + // Only the two COMPLETE rows survive, and totals reflect the FILTERED set (not 4). + assertEquals(2, resp.getTotalUniqueColumns()); + assertEquals(2, resp.getColumns().size()); + assertTrue( + resp.getColumns().stream().allMatch(c -> c.getMetadataStatus() == MetadataStatus.COMPLETE)); + assertNull(resp.getCursor()); + } + + @Test + void paginateFilteredItems_paginatesConsistentlyAcrossPages() { + List all = new ArrayList<>(); + for (int i = 0; i < 5; i++) { + all.add(item(String.format("col_%02d", i), MetadataStatus.COMPLETE, false)); + } + + ColumnAggregationRequest page1 = request("COMPLETE"); + page1.setSize(2); + ColumnGridResponse r1 = ColumnAggregator.paginateFilteredItems(all, page1); + + assertEquals(5, r1.getTotalUniqueColumns()); + assertEquals(2, r1.getColumns().size()); + assertEquals("col_00", r1.getColumns().get(0).getColumnName()); + assertEquals("col_01", r1.getColumns().get(1).getColumnName()); + assertNotNull(r1.getCursor(), "more pages remain, so a cursor is returned"); + + ColumnAggregationRequest page2 = request("COMPLETE"); + page2.setSize(2); + page2.setCursor(r1.getCursor()); + ColumnGridResponse r2 = ColumnAggregator.paginateFilteredItems(all, page2); + + assertEquals(5, r2.getTotalUniqueColumns()); + assertEquals(2, r2.getColumns().size()); + assertEquals("col_02", r2.getColumns().get(0).getColumnName()); + assertEquals("col_03", r2.getColumns().get(1).getColumnName()); + assertNotNull(r2.getCursor()); + + ColumnAggregationRequest page3 = request("COMPLETE"); + page3.setSize(2); + page3.setCursor(r2.getCursor()); + ColumnGridResponse r3 = ColumnAggregator.paginateFilteredItems(all, page3); + + assertEquals(1, r3.getColumns().size(), "last page has the remaining single item"); + assertEquals("col_04", r3.getColumns().get(0).getColumnName()); + assertNull(r3.getCursor(), "cursor is null on the last page"); + } } From 9d3caeb62fff5ea73f5f53170ad8c8e90a7e189d Mon Sep 17 00:00:00 2001 From: sonika-shah <58761340+sonika-shah@users.noreply.github.com> Date: Tue, 8 Sep 2026 08:53:26 +0530 Subject: [PATCH 2/4] perf(search): back column-grid status filter with a field-filtered _source scan MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the row-filter path's names-agg + per-name top_hits fan-out (~1+N queries/page) with a single _source scan per field-path group — the same mechanism the tag/glossary filter already uses — restricting _source to the column tree and entity-identity fields. This cuts queries per page from ~1+N to ~1 per field-path group, reuses the already-verified in-memory pagination, and reads every occurrence of a column instead of a 100-doc top_hits sample, so the aggregate status can no longer be misclassified (e.g. COMPLETE vs INCONSISTENT) by under-sampling. Restricting _source to columns + identity fields keeps the per-entity payload small, dropping the heavy entity-level derived fields (columnNames/columnNamesFuzzy). extractMatchingColumnsFromHit gains an includeAllColumns flag (the tag path passes false; the status scan passes true); the now-unused name-enumeration constants are removed. --- .../service/search/ColumnAggregator.java | 9 --- .../ElasticSearchColumnAggregator.java | 74 +++++++++++++------ .../OpenSearchColumnAggregator.java | 67 ++++++++++++----- 3 files changed, 100 insertions(+), 50 deletions(-) diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java index 5106dafd549d..14ea81ff4938 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java @@ -34,15 +34,6 @@ public interface ColumnAggregator { /** Max column names to retrieve in the names-only query during pattern search. */ int MAX_PATTERN_SEARCH_NAMES = 10000; - /** Terms include array matching every column name (row-level-filter path has no name pattern). */ - String MATCH_ALL_NAMES_REGEX = ".*"; - - /** - * Batch size for fetching occurrence data by column name in the row-level-filter path. Bounds the - * terms-include array and top_hits fan-out per query when many names must be materialized. - */ - int ROW_FILTER_DATA_BATCH = 1000; - /** * Number of sample docs pulled per column-name bucket to populate occurrences. Caps * {@code ColumnGridItem.totalOccurrences}; columns appearing in more entities than this diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java index b1b2d63fee42..6a15962ccb8a 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java @@ -282,20 +282,18 @@ private ColumnGridResponse aggregateColumnsWithPattern( /** * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The * filter acts on the aggregate status of a grouped column, which is only known after grouping all - * of its occurrences — so we enumerate every candidate name (respecting any column-name pattern - * and scope filters), fetch their occurrences, group them, then filter + paginate the items in - * memory. This keeps the page count and per-page size consistent with the filtered result set. + * of a column's occurrences. We read {@code _source} for the scoped entities in one scan per + * field-path group (the same mechanism the tag path uses), restricting {@code _source} to the + * column and identity fields, then group every column and filter + paginate the items in memory. + * This keeps the page count and per-page size consistent with the filtered result set, reads all + * occurrences (no top_hits sampling gap), and avoids one query per name. */ private ColumnGridResponse aggregateColumnsWithRowFilters( ColumnAggregationRequest request, List entityTypes) throws IOException { Map> fieldPathToEntityTypes = groupByFieldPath(entityTypes); - String regex = - !nullOrEmpty(request.getColumnNamePattern()) - ? ColumnAggregator.toCaseInsensitiveRegex(request.getColumnNamePattern()) - : MATCH_ALL_NAMES_REGEX; - - Map> allColumnsByName = new HashMap<>(); + Map> allColumnsByName = + new TreeMap<>(String.CASE_INSENSITIVE_ORDER); for (Map.Entry> entry : fieldPathToEntityTypes.entrySet()) { String columnNameKeyword = entry.getKey(); @@ -304,17 +302,7 @@ private ColumnGridResponse aggregateColumnsWithRowFilters( Query query = buildFilters(request, columnNameKeyword, null); try { - List names = executeNamesQuery(query, indexes, columnNameKeyword, regex).names(); - for (int i = 0; i < names.size(); i += ROW_FILTER_DATA_BATCH) { - List batch = names.subList(i, Math.min(i + ROW_FILTER_DATA_BATCH, names.size())); - Map> columnsByName = - executePageDataQuery(query, indexes, columnNameKeyword, columnFieldPath, batch); - for (Map.Entry> colEntry : columnsByName.entrySet()) { - allColumnsByName - .computeIfAbsent(colEntry.getKey(), k -> new ArrayList<>()) - .addAll(colEntry.getValue()); - } - } + fetchColumnsFromSource(indexes, query, columnFieldPath, allColumnsByName); } catch (ElasticsearchException e) { if (!isIndexNotFoundException(e)) { logShardFailureDetails(e, indexes, query); @@ -327,6 +315,46 @@ private ColumnGridResponse aggregateColumnsWithRowFilters( return ColumnAggregator.paginateFilteredItems(gridItems, request); } + /** {@code _source} fields needed to build grid items: entity identity + the whole column tree. */ + private List statusScanSourceIncludes(String columnFieldPath) { + return List.of( + "fullyQualifiedName", + "entityType", + "displayName", + "service.name", + "database.name", + "databaseSchema.name", + columnFieldPath); + } + + private void fetchColumnsFromSource( + List indexes, + Query query, + String columnFieldPath, + Map> columnsByName) + throws IOException { + + List includes = statusScanSourceIncludes(columnFieldPath); + SearchRequest searchRequest = + SearchRequest.of( + s -> + s.index(indexes) + .query(query) + .source(src -> src.filter(f -> f.includes(includes))) + .size(10000)); + + SearchResponse response = client.search(searchRequest, JsonData.class); + long totalHits = response.hits().total() != null ? response.hits().total().value() : 0; + if (totalHits > 10000) { + LOG.warn( + "Metadata-status source-fetch matched {} entities; only first 10000 scanned.", totalHits); + } + + for (Hit hit : response.hits().hits()) { + extractMatchingColumnsFromHit(hit, columnFieldPath, Set.of(), true, columnsByName); + } + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -427,7 +455,7 @@ private void fetchColumnsWithTagsFromSource( } for (Hit hit : response.hits().hits()) { - extractMatchingColumnsFromHit(hit, columnFieldPath, targetTags, columnsByName); + extractMatchingColumnsFromHit(hit, columnFieldPath, targetTags, false, columnsByName); } } @@ -435,6 +463,7 @@ private void extractMatchingColumnsFromHit( Hit hit, String columnFieldPath, Set targetTags, + boolean includeAllColumns, Map> columnsByName) { if (hit.source() == null) { return; @@ -457,7 +486,8 @@ private void extractMatchingColumnsFromHit( if (columnsData != null && columnsData.isArray()) { for (JsonNode columnData : columnsData) { String colName = getTextField(columnData, "name"); - if (colName != null && columnHasTargetTag(columnData, targetTags)) { + if (colName != null + && (includeAllColumns || columnHasTargetTag(columnData, targetTags))) { Column column = parseColumn(columnData, entityFQN); columnsByName .computeIfAbsent(colName, k -> new ArrayList<>()) diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java index 1c348bf1361d..2497e649fd02 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java @@ -211,25 +211,11 @@ private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationReque throws IOException { Query query = buildFilters(request, null); - String regex = - !nullOrEmpty(request.getColumnNamePattern()) - ? ColumnAggregator.toCaseInsensitiveRegex(request.getColumnNamePattern()) - : ColumnAggregator.MATCH_ALL_NAMES_REGEX; + Map> allColumnsByName = + new TreeMap<>(String.CASE_INSENSITIVE_ORDER); try { - List names = executeNamesQuery(query, regex).names(); - Map> allColumnsByName = new HashMap<>(); - for (int i = 0; i < names.size(); i += ColumnAggregator.ROW_FILTER_DATA_BATCH) { - List batch = - names.subList(i, Math.min(i + ColumnAggregator.ROW_FILTER_DATA_BATCH, names.size())); - Map> columnsByName = executePageDataQuery(query, batch); - for (Map.Entry> colEntry : columnsByName.entrySet()) { - allColumnsByName - .computeIfAbsent(colEntry.getKey(), k -> new ArrayList<>()) - .addAll(colEntry.getValue()); - } - } - + fetchColumnsFromSource(query, allColumnsByName); List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); return ColumnAggregator.paginateFilteredItems(gridItems, request); } catch (OpenSearchException e) { @@ -241,6 +227,47 @@ private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationReque } } + /** + * Read {@code _source} for the scoped entities in one scan, restricted to the column and identity + * fields, and extract every column. Backs the row-level-filter path (status / conflicts / missing + * metadata), which must group all of a column's occurrences before it can filter on the aggregate + * status. + */ + private void fetchColumnsFromSource( + Query query, Map> columnsByName) throws IOException { + + List resolvedIndexes = resolveIndexNames(); + List includes = + List.of( + "fullyQualifiedName", + "entityType", + "displayName", + "service.name", + "database.name", + "databaseSchema.name", + "columns"); + + SearchRequest searchRequest = + SearchRequest.of( + s -> + s.index(resolvedIndexes) + .query(query) + .source(src -> src.filter(f -> f.includes(includes))) + .size(10000)); + + SearchResponse response = client.search(searchRequest, JsonData.class); + long totalHits = response.hits().total() != null ? response.hits().total().value() : 0; + if (totalHits > 10000) { + LOG.warn( + "Metadata-status source-fetch matched {} entities; only first 10000 scanned.", totalHits); + } + + for (os.org.opensearch.client.opensearch.core.search.Hit hit : + response.hits().hits()) { + extractMatchingColumnsFromHit(hit, Set.of(), true, columnsByName); + } + } + /** * Tag/glossary filter path: the tag-check pass already extracted full column metadata from * _source (only tagged columns are in the map). Just paginate over the in-memory result. @@ -337,13 +364,14 @@ private void fetchColumnsWithTagsFromSource( for (os.org.opensearch.client.opensearch.core.search.Hit hit : response.hits().hits()) { - extractMatchingColumnsFromHit(hit, targetTags, columnsByName); + extractMatchingColumnsFromHit(hit, targetTags, false, columnsByName); } } private void extractMatchingColumnsFromHit( os.org.opensearch.client.opensearch.core.search.Hit hit, Set targetTags, + boolean includeAllColumns, Map> columnsByName) { if (hit.source() == null) { return; @@ -366,7 +394,8 @@ private void extractMatchingColumnsFromHit( if (columnsData != null && columnsData.isArray()) { for (JsonNode columnData : columnsData) { String colName = getTextField(columnData, "name"); - if (colName != null && columnHasTargetTag(columnData, targetTags)) { + if (colName != null + && (includeAllColumns || columnHasTargetTag(columnData, targetTags))) { Column column = parseColumn(columnData, entityFQN); columnsByName .computeIfAbsent(colName, k -> new ArrayList<>()) From 9f6d240bc6d9248b3c276577e114ffdb5fca88a3 Mon Sep 17 00:00:00 2001 From: sonika-shah <58761340+sonika-shah@users.noreply.github.com> Date: Tue, 8 Sep 2026 14:20:57 +0530 Subject: [PATCH 3/4] fix(search): honor columnNamePattern per column in the status-filter scan MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The _source-scan row-filter path returned every column of each matching entity (includeAllColumns=true). Combined with a columnNamePattern the entity-level name wildcard only decides which entities are scanned — flat-object mapping can't isolate the matching column — so non-matching columns leaked into the grid and inflated totalUniqueColumns. Drop columns whose name doesn't contain the pattern after the scan (mirrors the tag path), in both ES and OS aggregators. Also correct the row-filter Javadocs: OS still described the old name-enumeration path, and both overstated "reads all occurrences" — it reads every occurrence of each scanned entity, up to the 10K-entity scan cap. Adds an IT for the columnNamePattern + metadataStatus combination. --- .../it/tests/ColumnGridResourceIT.java | 45 +++++++++++++++++++ .../ElasticSearchColumnAggregator.java | 12 ++++- .../OpenSearchColumnAggregator.java | 19 ++++++-- 3 files changed, 71 insertions(+), 5 deletions(-) diff --git a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java index c4790fc4aee1..48203438108d 100644 --- a/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java +++ b/openmetadata-integration-tests/src/test/java/org/openmetadata/it/tests/ColumnGridResourceIT.java @@ -487,6 +487,51 @@ void test_getColumnGrid_metadataStatusPaginationCountsAreConsistent(TestNamespac assertAllRowsHaveStatus(page2, MetadataStatus.COMPLETE); } + @Test + void test_getColumnGrid_metadataStatusWithColumnNamePattern(TestNamespace ns) throws Exception { + OpenMetadataClient client = SdkClients.adminClient(); + DatabaseService service = DatabaseServiceTestFactory.createPostgres(ns); + DatabaseSchema schema = DatabaseSchemaTestFactory.createSimple(ns, service); + + // Two COMPLETE columns; only one name contains "alpha". + String matchName = ns.prefix("alpha_amount"); + String otherName = ns.prefix("zzz_other"); + for (String colName : List.of(matchName, otherName)) { + Column col = + Columns.build(colName) + .withType(ColumnDataType.BIGINT) + .withDescription("has description") + .withTags(List.of(new TagLabel().withTagFQN("PII.Sensitive"))) + .create(); + Tables.create() + .name(ns.prefix("pat_" + colName)) + .inSchema(schema.getFullyQualifiedName()) + .withColumns(List.of(col)) + .execute(); + } + waitForSearchIndexRefresh(ns); + + ColumnGridResponse response = + getColumnGrid( + client, + "entityTypes=table&metadataStatus=COMPLETE&columnNamePattern=alpha&serviceName=" + + service.getName()); + + // Combining columnNamePattern with a status filter must honor the pattern per column: + // the non-matching "zzz_other" column must not leak in (regression for the _source-scan path). + assertNotNull(response); + assertFalse(response.getColumns().isEmpty(), "the matching COMPLETE column should be returned"); + assertAllRowsHaveStatus(response, MetadataStatus.COMPLETE); + assertTrue( + response.getColumns().stream() + .allMatch(c -> c.getColumnName().toLowerCase().contains("alpha")), + "only columns whose name matches the pattern should be returned"); + assertEquals( + response.getColumns().size(), + response.getTotalUniqueColumns(), + "totalUniqueColumns must not include pattern-mismatched columns"); + } + /** * Every returned row must carry the requested aggregate status — the core guarantee of #26824 * (before the fix a COMPLETE/INCOMPLETE filter leaked rows of other statuses). diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java index 6a15962ccb8a..f04f8968ae3d 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java @@ -285,8 +285,8 @@ private ColumnGridResponse aggregateColumnsWithPattern( * of a column's occurrences. We read {@code _source} for the scoped entities in one scan per * field-path group (the same mechanism the tag path uses), restricting {@code _source} to the * column and identity fields, then group every column and filter + paginate the items in memory. - * This keeps the page count and per-page size consistent with the filtered result set, reads all - * occurrences (no top_hits sampling gap), and avoids one query per name. + * This keeps the page count and per-page size consistent with the filtered result set and reads + * every occurrence of each scanned entity, up to the {@code size(10000)}-entity scan cap. */ private ColumnGridResponse aggregateColumnsWithRowFilters( ColumnAggregationRequest request, List entityTypes) throws IOException { @@ -311,6 +311,14 @@ private ColumnGridResponse aggregateColumnsWithRowFilters( } } + // The name-pattern wildcard in buildFilters only decides which entities are scanned; + // flat-object + // mapping can't isolate the matching column. Drop columns whose name doesn't match, per column. + if (!nullOrEmpty(request.getColumnNamePattern())) { + String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); + allColumnsByName.keySet().removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); + } + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); return ColumnAggregator.paginateFilteredItems(gridItems, request); } diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java index 2497e649fd02..867294793143 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java @@ -203,9 +203,11 @@ private ColumnGridResponse aggregateColumnsWithPattern(ColumnAggregationRequest /** * Row-level-filter path (metadataStatus / hasConflicts / hasMissingMetadata, no tag filter). The * filter acts on the aggregate status of a grouped column, which is only known after grouping all - * of its occurrences — so we enumerate every candidate name (respecting any column-name pattern - * and scope filters), fetch their occurrences, group them, then filter + paginate the items in - * memory. This keeps the page count and per-page size consistent with the filtered result set. + * of a column's occurrences. We read {@code _source} for the scoped entities in one scan (the same + * mechanism the tag path uses), restricting {@code _source} to the column and identity fields, + * then group every column and filter + paginate the items in memory. This keeps the page count and + * per-page size consistent with the filtered set and reads every occurrence of each scanned + * entity, up to the {@code size(10000)}-entity scan cap. */ private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationRequest request) throws IOException { @@ -216,6 +218,17 @@ private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationReque try { fetchColumnsFromSource(query, allColumnsByName); + + // The name-pattern wildcard in buildFilters only decides which entities are scanned; + // flat-object mapping can't isolate the matching column. Drop non-matching columns per + // column. + if (!nullOrEmpty(request.getColumnNamePattern())) { + String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); + allColumnsByName + .keySet() + .removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); + } + List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); return ColumnAggregator.paginateFilteredItems(gridItems, request); } catch (OpenSearchException e) { From 58a4bcbbb1891765cfa0823d4029707840371d06 Mon Sep 17 00:00:00 2001 From: sonika-shah <58761340+sonika-shah@users.noreply.github.com> Date: Tue, 8 Sep 2026 14:31:49 +0530 Subject: [PATCH 4/4] refactor(search): extract shared applyColumnNamePattern helper The per-column name-pattern drop-loop was duplicated verbatim in the ES and OS row-filter scans. Move it to a static ColumnAggregator.applyColumnNamePattern helper (alongside the other shared row-filter helpers) so the two backends can't drift in pattern semantics. --- .../service/search/ColumnAggregator.java | 16 ++++++++++++++++ .../ElasticSearchColumnAggregator.java | 8 +------- .../opensearch/OpenSearchColumnAggregator.java | 11 +---------- 3 files changed, 18 insertions(+), 17 deletions(-) diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java index 14ea81ff4938..ab2ec73f3eee 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/ColumnAggregator.java @@ -20,6 +20,7 @@ import java.util.Base64; import java.util.Comparator; import java.util.List; +import java.util.Locale; import java.util.Map; import org.openmetadata.schema.api.data.ColumnGridItem; import org.openmetadata.schema.api.data.ColumnGridResponse; @@ -147,6 +148,21 @@ static boolean matchesRowFilters(ColumnGridItem item, ColumnAggregationRequest r return true; } + /** + * Drop columns whose name doesn't contain the request's {@code columnNamePattern} + * (case-insensitive). The name wildcard in the search query only scopes which entities are + * scanned; flat-object mapping can't isolate the matching column, so the pattern is enforced per + * column here. Shared by the ES and OS row-filter scans so their pattern semantics can't drift. + */ + static void applyColumnNamePattern( + Map columnsByName, ColumnAggregationRequest request) { + if (nullOrEmptyStr(request.getColumnNamePattern())) { + return; + } + String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); + columnsByName.keySet().removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); + } + /** A column has missing metadata if any of its groups lacks a description or tags. */ static boolean itemHasMissingMetadata(ColumnGridItem item) { if (item.getGroups() == null) { diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java index f04f8968ae3d..04c8d3e82e60 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/elasticsearch/ElasticSearchColumnAggregator.java @@ -311,13 +311,7 @@ private ColumnGridResponse aggregateColumnsWithRowFilters( } } - // The name-pattern wildcard in buildFilters only decides which entities are scanned; - // flat-object - // mapping can't isolate the matching column. Drop columns whose name doesn't match, per column. - if (!nullOrEmpty(request.getColumnNamePattern())) { - String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); - allColumnsByName.keySet().removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); - } + ColumnAggregator.applyColumnNamePattern(allColumnsByName, request); List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); return ColumnAggregator.paginateFilteredItems(gridItems, request); diff --git a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java index 867294793143..1442e02f2dc5 100644 --- a/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java +++ b/openmetadata-service/src/main/java/org/openmetadata/service/search/opensearch/OpenSearchColumnAggregator.java @@ -218,16 +218,7 @@ private ColumnGridResponse aggregateColumnsWithRowFilters(ColumnAggregationReque try { fetchColumnsFromSource(query, allColumnsByName); - - // The name-pattern wildcard in buildFilters only decides which entities are scanned; - // flat-object mapping can't isolate the matching column. Drop non-matching columns per - // column. - if (!nullOrEmpty(request.getColumnNamePattern())) { - String pattern = request.getColumnNamePattern().toLowerCase(Locale.ROOT); - allColumnsByName - .keySet() - .removeIf(name -> !name.toLowerCase(Locale.ROOT).contains(pattern)); - } + ColumnAggregator.applyColumnNamePattern(allColumnsByName, request); List gridItems = ColumnMetadataGrouper.groupColumns(allColumnsByName); return ColumnAggregator.paginateFilteredItems(gridItems, request);