diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index cfd0f4cfc..4ce5b1a31 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -1577,6 +1577,43 @@ removed: because ClickHouse compares a tuple ordering argument through a generic `Field` once per row per aggregate. +The native reader resolves that aggregate for **its own keys only**. An inner query +groups just the five sort-key columns under the page's filters and `LIMIT`, and +the outer query computes the `argMax` tuple for those keys. With a single +`GROUP BY ... LIMIT`, ClickHouse built the full 27-column tuple for every group +past the cursor before `LIMIT` kept `limit + 1` of them. The keyset comparison +`(tenant_id, experiment_id, run_id, captured_at_ns, capture_id) > (...)` is not +usable by the primary-key index, so each page cost about the whole catalog +whatever its size. Measured on 25.12 over a 198k-row catalog, a 38-row page +took ~150 ms, three quarters of it that tuple. Both queries carry every filter, +so the groups and their resolution are unchanged: that is the same immutability +rule, below, that makes the pre-aggregation filters safe. The Python reference +reader keeps the single-phase shape, so the parity suite compares the two +shapes directly. Each page still scans the rows past the cursor. Removing that +needs index-usable cursor bounds. + +The two queries must filter alike. If the inner query drops a filter that the +outer one keeps, such as snapshot membership or a hook filter, its `LIMIT` fills +with keys the outer query then discards. The page comes back short, a short page +issues no cursor, and the walk ends early while reporting success. The parity +suite walks a corpus built for this: an unpublished pack whose keys interleave +with the member keys, crossed in pages of 2, with and without a hook filter. + +The second read counts against the read guard. `max_rows_to_read` limits the +whole statement, and both queries read `*_capture_raw`. The inner query reads +every row past the cursor that the filters do not prune. The outer query reads +every granule that holds one of the page's keys. A selective filter spreads +those keys across granules, so the outer read can approach a second full scan, +and a page can read up to **twice** the rows the single-phase query did. On one +corpus a page that read 1,102,131 rows now reads 2,203,414, and a +`max_rows_to_read` of 1,500,000 that the old shape passed refuses the new one +with Code 158 (`TOO_MANY_ROWS`). On the parity test's corpus, one granule per +part, 25.12 reports 68 rows read for a native page against 34 for the Python +reader's single-phase page. Size `ReaderConfig::max_rows_to_read` (`reader.h`) +for two passes over the rows past the cursor. At worst, the default of +50,000,000 then trips once about 25 million rows lie past the cursor, not 50 +million. + Every descriptor field except the locator is immutable for a `(tenant_id, capture_id)`, which is what makes the pre-aggregation `WHERE` filters safe; the rule is written out in `clickhouse_reader`'s module diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index a83ff32ae..7f2d3342f 100644 --- a/native/csrc/catalog/reader.cpp +++ b/native/csrc/catalog/reader.cpp @@ -766,9 +766,35 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { grouped += quoted(column); order += quoted(column); } + // Choose the page's keys first, then resolve the argMax tuple for those keys + // alone. One GROUP BY ... LIMIT built the full resolution tuple for EVERY + // group past the cursor before LIMIT kept limit + 1 of them, and the keyset + // tuple comparison is not usable by the primary-key index, so a page cost + // about the whole catalog whatever its size. The Python reference reader + // keeps that single-phase shape, and the parity suite compares the two. + // + // Both queries carry every filter (the same `clauses`), so groups and + // resolution are unchanged. An inner query missing one -- snapshot + // membership, a hook filter -- fills its LIMIT with keys the outer query + // then drops: the page comes back short, owes no cursor, and a walk ends + // early. test_a_page_walk_skips_unpublished_keys_without_ending_early pins + // that. + // + // The second read counts against the read guard. max_rows_to_read limits + // the whole statement, and both queries read capture_raw: the inner one + // every row past the cursor that the filters leave, the outer one each + // granule holding a page key. A selective filter spreads those keys across + // granules, so the outer read approaches a second full scan, and a page can + // read up to twice the rows the single-phase query did (1,102,131 -> + // 2,203,414 on one corpus, so a limit of 1,500,000 that passed the old + // shape refuses this one with Code 158). Size max_rows_to_read for two + // passes over the rows past the cursor. const std::vector rows = client_->execute( "SELECT " + projection() + " FROM " + qualified("capture_raw") + - " WHERE " + clauses + " GROUP BY " + grouped + " ORDER BY " + order + + " WHERE " + clauses + " AND (" + grouped + ") IN (SELECT " + + grouped + " FROM " + qualified("capture_raw") + " WHERE " + + clauses + " GROUP BY " + grouped + " ORDER BY " + order + + " LIMIT %(row_limit)s) GROUP BY " + grouped + " ORDER BY " + order + " LIMIT %(row_limit)s", params, bounded_read_settings()); diff --git a/tests/test_native_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 83d3e2aab..4c7f30883 100644 --- a/tests/test_native_reader_parity_live.py +++ b/tests/test_native_reader_parity_live.py @@ -2614,3 +2614,169 @@ def test_a_non_bmp_hook_name_indexes_intact(fake_s3): sink.close() store.close() driver.close() + + +def test_native_pages_resolve_argmax_for_their_own_keys_and_match_python(): + """The native page query resolves the argMax tuple for the page's keys only. + + A single GROUP BY ... LIMIT built the full 27-column resolution tuple for + every group past the cursor before LIMIT kept limit + 1 of them, and the + keyset tuple comparison is not usable by the primary-key index -- so each + page cost about the whole catalog whatever its size (~150 ms for a 38-row + page over 198k rows). The native reader now selects the page's keys in an + inner, key-only query under the same filters and LIMIT. + + Two things are pinned. The SHAPE, from the query ClickHouse actually + received: native reads arrive over HTTP (interface 2), and every argMax + query must carry the key subquery. And the RESULT across a superseded + version, walked in pages small enough to cross several cursors, against + the Python reader, which still uses the single-phase shape -- so the + parity suite checks the new shape against the old one. + """ + from dmi.storage.capture.clickhouse_reader import _RESOLVED + + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + first = _descriptor_dicts(10) + _publish_native(driver, prefix, first, 7) + # Re-describe four captures from a different pack at a later + # version: resolution must pick version 8's locator for those. + newer_pack = "018f0000-0000-7000-8000-00000000beef" + second = [dict(first[i], pack_id=newer_pack) for i in (0, 3, 6, 9)] + _publish_native(driver, prefix, second, 8) + + native, cursor = [], None + while True: + page = driver.call(op="search", limit=3, + **({"cursor": cursor} if cursor else {})) + assert page["ok"], page + native += page["items"] + cursor = page["next_cursor"] + if cursor is None: + break + reader = _python_reader(client, config) + python, cursor = [], None + while True: + page = _python_page_items(reader, limit=3, cursor=cursor) + python += page.items + cursor = page.next_cursor + if cursor is None: + break + + assert len(native) == 10 == len(python) + assert _normalize(native) == _normalize(python) + pack_at = 5 + list(_RESOLVED).index("pack_id") + moved = {row[4] for row in native if row[pack_at] == newer_pack} + assert moved == {"capture-0", "capture-3", "capture-6", "capture-9"} + + client.execute("SYSTEM FLUSH LOGS") + sent = [row[0] for row in client.execute( + "SELECT query FROM system.query_log WHERE type = 'QueryFinish' " + "AND interface = 2 AND query LIKE %(raw)s AND query LIKE '%%argMax%%' " + "AND event_time > now() - 600", + {"raw": f"%{prefix}_capture_raw%"})] + assert sent, "no native page query reached ClickHouse over HTTP" + # The KEY tuple feeds the subquery. A bare ") IN (SELECT " would also + # match the snapshot-membership filter `(store_id, pack_id) IN + # (SELECT ...)` that the old single-phase query already carried. + keys = "`tenant_id`,`experiment_id`,`run_id`,`captured_at_ns`,`capture_id`" + for query in sent: + assert f"({keys}) IN (SELECT {keys} FROM" in query, query + inner = query.split(f"({keys}) IN (SELECT {keys} FROM", 1)[1] + assert "argMax" not in inner.split(" GROUP BY ", 1)[0], query + assert query.count(" LIMIT ") == 2, query + finally: + driver.close() + + +def test_a_page_walk_skips_unpublished_keys_without_ending_early(): + """Every member capture comes back once, however unpublished keys interleave. + + The two-phase page is correct only while its inner key query and its outer + resolution query filter alike. Let the inner one drop a filter the outer one + keeps -- snapshot membership, a hook filter -- and its LIMIT fills with keys + the outer query then discards: the page comes back short, a short page owes + no cursor, and the walk stops early while reporting success. The other tests + cannot see that, because every row they stage is a member that matches the + filter, so dropping the inner membership filter left them all passing. + + Here a published pack holds ten captures whose hooks alternate, and an + unpublished pack -- written at version 8 and never published, as a crashed + indexer leaves one -- holds two keys between each pair of member keys, all + under the filtered hook. It also re-describes two members, and only + membership keeps that later version from winning their argMax. Walked in + pages of 2 (an inner LIMIT of 3, so the first page alone crosses two + unpublished keys), with and without a hook filter, the native reader must + return exactly the member captures, in order and resolved to the published + pack, and agree with the Python reader. + """ + from dmi.storage.capture.clickhouse_reader import _RESOLVED + + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + published_pack = "018f0000-0000-7000-8000-00000000a001" + unpublished_pack = "018f0000-0000-7000-8000-00000000b002" + # captured_at_ns rises with the index, so the index is sort order. + members, unpublished = [], [] + for index, entry in enumerate(_descriptor_dicts(30)): + if index % 3 == 0: + hook = "resid_pre" if index % 6 == 0 else "other_hook" + members.append( + dict(entry, pack_id=published_pack, hook_name=hook)) + else: + unpublished.append( + dict(entry, pack_id=unpublished_pack, + hook_name="resid_pre")) + unpublished += [dict(members[i], pack_id=unpublished_pack) + for i in (1, 4)] + _publish_native(driver, prefix, members, 7) + written = driver.call(op="write_descriptors", + descriptors=unpublished, index_version=8) + assert written["ok"], written + assert driver.call(op="current_watermark")["watermark"] == "7" + + reader = _python_reader(client, config) + pages = len(members) + len(unpublished) + 2 # a walk that loops fails + + def walk_native(**filters): + rows, cursor = [], None + for _ in range(pages): + page = driver.call( + op="search", limit=2, **filters, + **({"cursor": cursor} if cursor else {})) + assert page["ok"], page + rows += page["items"] + cursor = page["next_cursor"] + if cursor is None: + return rows + raise AssertionError(f"native walk did not end: {filters}") + + def walk_python(**filters): + items, cursor = [], None + for _ in range(pages): + page = _python_page_items( + reader, limit=2, cursor=cursor, **filters) + items += page.items + cursor = page.next_cursor + if cursor is None: + return items + raise AssertionError(f"python walk did not end: {filters}") + + pack_at = 5 + list(_RESOLVED).index("pack_id") + for native_filters, python_filters, expected in ( + ({}, {}, [m["capture_id"] for m in members]), + ({"hook_names": ["resid_pre"]}, {"hook_names": ("resid_pre",)}, + [m["capture_id"] for m in members + if m["hook_name"] == "resid_pre"]), + ): + native = walk_native(**native_filters) + assert [row[4] for row in native] == expected, native_filters + assert {row[pack_at] for row in native} == {published_pack} + python = walk_python(**python_filters) + assert _normalize(native) == _normalize(python), native_filters + finally: + driver.close()