The problem isn't too dissimilar from the challenge that tracing profilers try to solve. One key difference, of course, is that dynamically adapting the physical storage strategy to usage patterns is relative expensive, resource-wise. But not every problem needs to be solved for those users with 1 billion rows. I would argue that most applications — web apps, anyway — have datasets that are large enough to optimize, yet still small enough that a dynamic, self-tuning system could be entirely feasible.
The resource problem can be mitigated in several ways. One is to off-load processing horizontally, so while a master is chugging away, a slave is building a table using a different strategy. Another is to simply perform expensive optimizatons during off-peak hours.
Also, I find it interesting that databases also still operate in a mode where reads and writes happen in the same physical store. I have an idea for a database that separates the two: You have the transactional, highly consistent, non-horizontally-scalable OLTP-style "transaction store" where a client carefully constructs reads and writes against the absolute truth, but where reads aren't super efficient; and then you have the "query store" which is a read-only, highly replicated and sharded set of mirrors that uses highly optimized data structures that are created on demand based on the topology of the data and usage patterns (in particular, the same data can potentially be indexed in multiple ways (B-tree, compressed columns, etc.) until the profiler finds the best strategy). The mirrors' replication should be tunable so that some "hot" subsets can be immediately consistent, and others can use slower batch transfers to reduce lock contention on updates.
Right now, when people build scalable systems, they often apply this design by putting the truth in PostgreSQL or whatever, and then indexing it into Elasticsearch or Solr. But there are huge benefits to combining the two into a single, unified system.