3 ms·
The sharding thing is such an interesting problem. I'm thinking about the case where a user searches for "A B". Assume A has a few million matching IDs, B als
by slackerIII 19y ago
The sharding thing is such an interesting problem. I'm thinking about the case where a user searches for "A B". Assume A has a few million matching IDs, B also has few million, and (of course) the indexes for A & B are stored on different servers. Do you have to pull the full result sets for A & B back to a single aggregator to intersect them? I'd love to learn more about the optimizations for that problem.
- elq 19y agoAh. Well, keep in mind that we're splitting the documents up in an arbitrary way WRT the data, so every cluster/partition has a complete copy of the index for the documents that happen to be bucketed there. And no cluster owns a single "document class/cluster". So, if you want to search for, say "microsoft windows", the aggregator just sends something like "PHRASE('microsoft', 'windows')", each query node finds the document vector/set for microsoft and the document vector for windows and does an intersection of the doc ids, then the node has to do scan that set, grab the document position array from each document and filters out any documents where windows doesn't occur at microsoft-Positions + 1. All of the conjunctions, disjunctions, wildcard expansions, near operations, and phrase operations, etc are executed on the query node. All of the complex sort evaluation also happen on the query node. The aggregator only merges result sets and and performs any necessary global sorting.
- slackerIII 19y agoAh, of course. I misunderstood your previous comment. The system is partitioned by document, not by token - that makes a lot more sense. Oh well, I guess that isn't quite as hard of a problem as I thought :) Thanks for the follow up.