Skip to content

fix: merge multi-shard reads in the key sorting of the namespace - #1388

Merged
merlimat merged 3 commits into
oxia-db:mainfrom
merlimat:fix-client-multi-shard-key-sorting
Sep 30, 2026
Merged

merlimat merged 3 commits into
oxia-db:mainfrom
merlimat:fix-client-multi-shard-key-sorting

Conversation

@merlimat

Copy link
Copy Markdown
Collaborator

Problem

Some reads without a partition key go to every shard of the namespace:

  • a Get with ComparisonFloor, ComparisonLower, ComparisonCeiling or ComparisonHigher, or with UseIndex;
  • a RangeScan.

The client then has to combine the shards' answers. For a get it picks one answer: the largest for floor/lower, the smallest for ceiling/higher. For a range scan it merges the shards' streams into one ordered stream.

Each shard orders its keys by the key sorting of the namespace:

  • natural: by their bytes;
  • hierarchical (the default): keys with fewer / first, then by their bytes, with / after any other byte.

The client compares keys with compare.CompareWithSlash instead. That is the order of the old storage format, used before v0.15, and it matches neither sorting.

So a floor or ceiling get can return the wrong record, and a range scan can return records out of order, with both sortings. The client can't fix this on its own, because only the namespace config holds the key sorting. The shard assignments the client receives don't carry it.

Example

A natural namespace with 4 shards:

  1. Put "a0" (it lands on shard 2) and "a/y/z" (it lands on shard 0).
  2. Call Get("b", ComparisonFloor()), which asks for the greatest key ≤ "b". In byte order "a/y/z" < "a0" < "b" (/ is 0x2F, 0 is 0x30), so the answer should be "a0".
  3. The get has no partition key, so the client sends it to all 4 shards. Shard 2 answers "a0", shard 0 answers "a/y/z", and the other two answer "not found".
  4. Before: the client keeps the larger answer according to CompareWithSlash. That order puts a key with a / after one without, so it ranks "a0" < "a/y/z". The get returns "a/y/z", a key lower than "a0".
  5. After: the client received KEY_SORTING_NATURAL with the shard assignments. It compares the bytes and returns "a0".

The default sorting is affected too. Take a hierarchical namespace with "a/x" (shard 3) and "ab/y" (shard 0), and call Get("a/z", ComparisonFloor()). In hierarchical order "ab/y" < "a/x" < "a/z", since / sorts after b, so the answer should be "a/x".

  • Before: it returns "ab/y". CompareWithSlash compares the first segments, "a" < "ab".
  • After: it returns "a/x".

A RangeScan("", "") over b, a/y/z, ab/y, a0, a/x on 4 shards returns:

before after
natural a0, b, a/x, a/y/z, ab/y a/x, a/y/z, a0, ab/y, b
hierarchical a0, b, a/x, ab/y, a/y/z a0, b, ab/y, a/x, a/y/z

Modification

  • Protocol: a new enum KeySorting (KEY_SORTING_UNKNOWN, KEY_SORTING_NATURAL, KEY_SORTING_HIERARCHICAL) and a new field NamespaceShardsAssignment.key_sorting (field 3). It is a new enum rather than replication.KeySortingType because client.proto can't import replication.proto.

  • Coordinator: fills the field for each namespace from the namespace config. A namespace with no key sorting is sent as hierarchical, which is what its shards use. The standalone server fills it from its own config.

  • Data servers: unchanged. They pass the assignments on to the clients. Older data servers (v0.15.0, v0.16.1 and v0.17.1 checked) keep the field as an unknown field, so it still reaches the clients.

  • Go client:

    • The shard manager keeps the key sorting from the latest assignments.
    • Multi-shard gets compare keys in that sorting, using the same key encoder as the shards. Index gets compare the secondary keys first.
    • The range-scan merge uses the same order, and computes each record's comparison key once.
  • Older coordinators: when the coordinator doesn't send the field, the client keeps comparing with CompareWithSlash, exactly as before. This way a client-only upgrade doesn't change any result. It also stays correct for servers older than v0.15, which store keys in that order.

  • Compatibility: old clients ignore the field. The fix takes effect once the coordinator (or the standalone server) and the client are upgraded.

  • Tests: the example is reproduced by these e2e tests, which fail on main:

    • TestSyncClientImpl_FloorCeilingGetKeySorting (10 get cases);
    • TestSyncClientImpl_RangeScanKeySorting;
    • TestCoordinator_MultiShardKeySorting, which goes through a real coordinator and data server.

    With a single shard, where the server alone decides, all their cases pass, so the expected values are right. Unit tests cover:

    • the key order;
    • index gets;
    • the range-scan merge, including the fallback;
    • the shard manager;
    • the coordinator's assignments.

Notes

  • Not fixed here: a range scan with UseIndex is still merged by the primary key, because the range-scan records don't carry their secondary key. That needs a server change, left as a follow-up.
  • Index gets: they are correct end to end only together with fix: order secondary-index gets by the shard's key sorting #1387, which fixes the server-side comparison.
  • Other SDKs: the Java and Rust SDKs have the same bug. Their changes are described separately and are not part of this PR.
  • API change: oxia.ResultHeap, exported by accident, changes from a slice to a struct.

A floor, lower, ceiling or higher get without a partition key goes to every
shard, and the client picks one of the answers with compare.CompareWithSlash.
That order matches neither key sorting of the shards, so the client can pick
the wrong answer. The test stores two keys on different shards and checks the
answer for natural and hierarchical namespaces. It fails on main.

Signed-off-by: Matteo Merli <mmerli@apache.org>
A range scan without a partition key goes to every shard, and the client
merges the records of the shards with compare.CompareWithSlash. That order
matches neither key sorting of the shards, so the merged records are out of
order. The test scans keys stored on several shards for natural and
hierarchical namespaces. It fails on main.

Signed-off-by: Matteo Merli <mmerli@apache.org>
A get with a floor, lower, ceiling or higher comparison or with an index, and a
range scan, go to every shard when they have no partition key, and the client
picks or merges the results of the shards. The client compared the keys with
compare.CompareWithSlash, the order of the old storage format. It matches
neither key sorting of the shards, so the client could return the wrong record
or the records out of order. For example, in a natural namespace with "a0" and
"a/y/z" on different shards, a floor get of "b" returned "a/y/z" instead of
"a0"; in a hierarchical namespace with "a/x" and "ab/y", a floor get of "a/z"
returned "ab/y" instead of "a/x".

The clients can't find the order by themselves: only the namespace config has
the key sorting. The coordinator, and the standalone server, now send it with
the shard assignments of each namespace, in the new key_sorting field of
NamespaceShardsAssignment. A namespace without a key sorting is sent as
hierarchical, the order that the data servers use for it. The data servers
relay the field with the rest of the assignments; older data servers keep it
as an unknown field, so it reaches the clients through them too.

The Go client compares the keys, and the secondary keys of an index get, with
the encoder of that key sorting. When the server does not send the key
sorting, the client keeps comparing the keys as before.

Signed-off-by: Matteo Merli <mmerli@apache.org>
Copilot AI balanced review requested due to automatic review settings September 29, 2026 03:20

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@merlimat
merlimat merged commit 0217793 into oxia-db:main Sep 30, 2026
5 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants