3 ms·
This is interesting. In effect, you're doing your sharding client-side? What do you plan to do if you need to add new nodes and/or you need to rebalance nodes
by tepidandroid 7y ago
This is interesting. In effect, you're doing your sharding client-side?
What do you plan to do if you need to add new nodes and/or you need to rebalance nodes due to concentrated data access patterns? How do you handle cross-node queries like joins?
- bsg75 7y ago> In effect, you're doing your sharding client-side? Correct. Some ETL inserts to specific nodes when sharding is necessary, in other cases Kafka engine tables on a group of nodes subscribe to common topics, and we simply let the whole cluster participate in queries. This works just fine when table scans are acceptable. Rebalancing is a missing option here, short of moving partitions manually. But in my specific use cases, I have not yet needed to rebalance across nodes. Note using native Clickhouse replication is still an option if we need it. One cost to it is the extra work needed in the database cluster, so addressing it in an eariler layer works for us. > How do you handle cross-node queries like joins? If I understand your question, since we are using the Distributed view type across the cluster definition, a query on any node will receive data from the others as part of a join-less SELECT, and federate on the node with the client connection. We are not doing any database-side JOINs currently. Plans are to augment data in ETL, or join data post-query (potentially Spark). Clickhouse dictionaries handle simple cases.
- tepidandroid 7y agoThanks for the details, this is very helpful.
- valyala 7y agoThere is a chproxy [1] - a proxy that is able to balance inserts and selects across ClickHouse nodes / replicas. Then client applications shouldn't know anything about ClickHouse cluster topology - they just talk to chproxy. [1] https://github.com/Vertamedia/chproxy https://github.com/Vertamedia/chproxy