If you're using Cassandra and you're also not a fool, it's probably because you actually need a lot of scale, so you'll generally do the denormalization up front and you will denormalize everything. If you're writing an event, you will write it several places. For example, suppose you're WalMart recording sales. You might write to: store transactions by store/year/day/hour (the "master" record insofar as you have one), user transactions by user/year/month, product purchases by manufacturer/year/day/hour... When you write the transaction, you write to all of these locations.
Each of these "by X" keys is a shard. Each can be located on a different set of ~3 machines (the number is configurable). Querying involves getting a copy of the ring topology, computing which integer shard-ID the key maps to, figuring out which machines in the ring own that integer, and then asking the machine for a whole bucketful of data, which should be a superset of what you're actually looking for. For something like a user's transactions you'll want to have basically everything there at once, so loading the "order history" page for the past month might be a single query that just returns a report: no joins at all, very fast, super scalable. Other lookup strategies might ask for a range within that bucket (the data within the bucket can be ordered by a single key; often this is a timestamp or time-based UUID). Anything that isn't a simple query of a few buckets like this has to be a map-reduce job and will be slow.
All of this is pain. You should generally not invite pain into your organization. However, if pain has already found you, something like this may be the least painful option.
Not really a superset, cassandra (and scylla, and bigtable, all of which are basically copying bigtable's model) each try very hard not to read any extra data at all, and can often return approximately the exact data requested, modulo serialization data which is usually fitting in ~compression chunk size (64k) + checksum.
> If you're using Cassandra and you're also not a fool, it's probably because you actually need a lot of scale
Cassandra also gives you very literally the most control over CAP tradeoffs of any database in the industry.
If you have 100 machines per DC in each of 10 dcs, what happens when one machine is offline? one rack? one dc? 2 dcs separated from 8? 6 dcs separated from 4? There's no single answer in cassandra (depends on replication factor, consistency of writes, consistency of reads, all of which are tunable, with 2 of those being tunable PER QUERY), the CAP tradeoffs are yours and yours alone. That flexibility is powerful for power users (it's also confusing for novices, which is unfortunate).
But to your first point, yes, the point is scale. The lack of opinions and deliberate functionality are designed to enable it to scale to thousands of hosts, potentially petabytes of data, trivially accessible in a single SQL-like CQL query, with realistic read latency < 1ms mean/avg and < 5ms p99 for a tuned workload where you know what you're doing. A lot of users will never need a database that can do a million reads per second across a thousand machines reaching p50 1ms on 2 petabytes of data, but Cassandra can do that, and you don't have to build a whole sharding layer on top of mysql/postgres/redis or even install Scylla to get there.