Repository navigation
fix: merge multi-shard reads in the key sorting of the namespace - #1388
Merged
merlimat merged 3 commits intoSep 30, 2026
Merged
Conversation
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>
merlimat
requested review from
RobertIndie,
coderzc and
mattisonchao
as code owners
September 29, 2026 03:20
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
Some reads without a partition key go to every shard of the namespace:
GetwithComparisonFloor,ComparisonLower,ComparisonCeilingorComparisonHigher, or withUseIndex;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:
/first, then by their bytes, with/after any other byte.The client compares keys with
compare.CompareWithSlashinstead. 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:
"a0"(it lands on shard 2) and"a/y/z"(it lands on shard 0).Get("b", ComparisonFloor()), which asks for the greatest key ≤"b". In byte order"a/y/z" < "a0" < "b"(/is 0x2F,0is 0x30), so the answer should be"a0"."a0", shard 0 answers"a/y/z", and the other two answer "not found".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".KEY_SORTING_NATURALwith 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 callGet("a/z", ComparisonFloor()). In hierarchical order"ab/y" < "a/x" < "a/z", since/sorts afterb, so the answer should be"a/x"."ab/y".CompareWithSlashcompares the first segments,"a" < "ab"."a/x".A
RangeScan("", "")overb, a/y/z, ab/y, a0, a/xon 4 shards returns:Modification
Protocol: a new enum
KeySorting(KEY_SORTING_UNKNOWN,KEY_SORTING_NATURAL,KEY_SORTING_HIERARCHICAL) and a new fieldNamespaceShardsAssignment.key_sorting(field 3). It is a new enum rather thanreplication.KeySortingTypebecauseclient.protocan't importreplication.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:
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:
Notes
UseIndexis 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.oxia.ResultHeap, exported by accident, changes from a slice to a struct.