3 ms·
Reposting my comment from when this came up on our 1.0 launch: IIRC, in the case that aphyr refers to for these specific numbers, the reads are scans that span
by arjunnarayan 9y ago
Reposting my comment from when this came up on our 1.0 launch:
IIRC, in the case that aphyr refers to for these specific numbers, the reads are scans that span multiple shards[1], while the writes are writes to single shards.
[1] even though aphyr says it's just a hundred rows, the tables are split into multiple shards because aphyr in this case was specifically testing our correctness in multi-shard contention scenarios. In production you wouldn't have multi-shard reads crop up until you were doing scans for tables that were hundreds of thousands of rows in size[2]. It's easier to picture if you think of this performance speed difference in the scenario where you are doing full table scans spanning multiple shards on multiple computers while the underlying rows are being rapidly mutated by contending write transactions. The transactions get constantly aborted to avoid giving a non-serializable result, and performance is suffering. We agree that the numbers in this contention scenario are too low, and we are actively working on high-contention performance (and performance in general) leading up to our 1.0 release[3].
[2] Specifically, we break up into multiple shards when a single shard exceeds 64mb in size.
[3] You can follow along on one of the PRs that address this specific performance issue here: https://github.com/cockroachdb/cockroach/pull/13501 https://github.com/cockroachdb/cockroach/pull/13501