Our data isn't enormous, by any means - 160G in one particular instance that's being used for proof of concept, it'll add up to 5-20T should it reach production. The catch is that it's 160G in MySQL; it's only 5G or so once it's been boiled down to Parquet files in HDFS. Columnar stores can be a really big win, depending on the shape of your data.
We use Impala for our queries. It's quite good tech; it's much faster at table scans than everything that doesn't describe itself as an in-memory database. That means writing SQL much like you would with Hive, only it runs faster.
I tried out both Citus and Greenplum to give PostgreSQL a fair shot. First problem: PostgreSQL is limited to 1600 columns in a table, and the column limit for a select clause isn't much bigger. We have several times this number of columns in our largest analytic tables. Not the end of the world, you can cobble things together from joins and more special-purpose tables.
Second problem: CitusDB doesn't come OOTB with a column store, and it's far too slow when using a row store. I didn't bother trying to compile the column store extension to use with Citus; the pain ruled itself out. I continued ahead with Greenplum, focusing on columnar storage - row storage is consistently poor.
Third problem: Greenplum is cobbled together from a pile of duct tape, an assembly of scripts and ssh keys to keep the cluster in sync. It does not inspire the same kind of confidence for operational management as HDFS, whether rebalancing the cluster, expanding the cluster, or decommissioning nodes (not supported with Greenplum, AFAICT).
Fourth problem: Impala simply runs faster than the Postgres derivatives, and its lead increases the more data you have. Impala seems to do table scans over twice as fast on identical deployment environments.
Indexes only help when the operation being performed can use the index. As it happens, most analytic queries do full scans, or have predicates that are either not very selective (randomly skipping rows here and there) or are really selective (date bucketing, which maps well to typical Hadoop partitioning strategies). I had some hope that indexes would help for joins; but Greenplum didn't elect to use my indexes, and when I forced their use, it ran slower. The ancient version of Postgres that Greenplum is forked from doesn't help much either, since it can't e.g. use covering indexes to avoid looking back to the table.
If it was my startup, I'd take a risk on something like MonetDB, or look harder at MemSQL, given what I've seen about how data has shrunk with column stores. But from what I've seen and measured, Postgres doesn't really cut it for analytic queries.