Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 37 additions & 0 deletions docs/capture-storage-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 27 additions & 1 deletion native/csrc/catalog/reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<Row> 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());

Expand Down
166 changes: 166 additions & 0 deletions tests/test_native_reader_parity_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Loading