Maybe I am wrong, but I am yet to read stories on operating tens of TB scale workloads on citus.
The biggest problem of Citus was migration effort, transition from single node to multi-node was not trivial. Here I’m not talking about single table use-cases, more classic relational, multi-tenant apps with 100s to 1000s of tables. This is partly expected with most sharding technologies, though.
Sharing some insights based on my multiple years of experience working with Citus!
Here are few customer use-cases I could found:
https://docs.citusdata.com/en/v10.0/get_started/what_is_citu...
https://info.citusdata.com/rs/235-CNE-301/images/Citus_Data_...?
Schema changes need locking... everywhere. Citus is no different. And they need proper design everywhere. If you mean that it requires distributed transactions, well, yes, again: expected and solved. Not even all sharding solutions support this.
> coordinator node
Not sure what the problem is here. If what you mean is that a single coordinator, even with an HA replica, can saturate, that's true, but you can add multiple "query routers" (that's our name in StackGres, see [1]).
> Maybe I am wrong, but I am yet to read stories on operating tens of TB scale workloads on citus.
For example, we have a customer that ingests some 30TB/day, and it's ramping up towards 200TB/day of ingestion. On 24 worker nodes.
[1]: https://stackgres.io/doc/latest/administration/sharded-clust...