Auxiliary Index Join Query Parser

auxIndexJoin (which stands for Auxiliary Index Join) Query Parser is similar to Join Query Parser, but uses a lazily written sidecar index for faster joins.

Instead of using Lucene’s join utilities, the parser matches documents through a dedicated join index that is maintained alongside the core’s main index. The join index is populated lazily: if the inner (from-side) query doesn’t hit a certain segment, the join index isn’t written for it; and if {!auxIndexJoin} is intersected (AND/+) with a query that doesn’t hit a certain segment on the outer (to) side, the corresponding join index columns aren’t written either.

Like the Join query parser, Solr runs a subquery (the v parameter), gathers the values that matching documents have in a from field, and returns documents where those values are contained in a to field.

For example:

q={!auxIndexJoin from=manu_id_s to=id}title:ipod

Parameters

This query parser takes the following parameters:

from

Required

Default: none

The "foreign key" field name, collected while enumerating the subordinate query. So far, it should be single value string field with docValues enabled.

to

Required

Default: none

The "primary key" field name looked up in the local core’s index. So far, it should be single value string field with docValues enabled.

fromIndex

Optional

Default: processing core

The name of the core to run the "from" query (v parameter) on and where "from" values are gathered. If this parameter is not defined, it defaults to the processing core. Cross-core joins are the primary use case for auxIndexJoin, so this is useful when the "from" values live in a different core on the same node.

Configuration

auxIndexJoin query parser must be registered in solrconfig.xml as a <queryParser>. A sidecar join index is opened per core when the core loads, in a directory under the core’s dataDir, and closed when the core closes.

dir

Optional

Default: aux-index-join

The init parameter for the directory holding the sidecar join index. Resolved relative to the core’s dataDir unless an absolute path is given.

singleFieldPerSegment

Optional

Default: false

If true, each pair column is flushed into its own sidecar segment. Otherwise, every pair column built in the same round is batched into a single segment (the default), trading a longer sweep for less storage overhead.

blockingRefresh

Optional

Default: true

If true, writing a batch of pair columns blocks until the sidecar’s searcher manager is refreshed past it, so the freshly built pairs are immediately visible to the caller that triggered the build.

useFromSideThreads

Optional

Default: true

If true (the default), searching the from-side segments and loading foreigh key columns is parallelized across the executor threads described in Segment-level Parallelism. Note it’s actual only if Segment-level Parallelism is enabled and treated as false otherwise. If false, the from-side segments are loaded sequentially on the calling thread, bounding the build to a single thread.

sweepSamplingInterval

Optional

Default: 60

How often (in seconds) the dead-pair reaper samples searcher state while the sidecar index is being read. Calls arriving sooner than this interval after the last accepted sample are skipped. Pass a non-positive value to sample on every call.

<queryParser name="auxIndexJoin" class="org.apache.solr.search.join.AuxIndexJoinQParserPlugin">
  <str name="dir">aijoin</str>
  <bool name="singleFieldPerSegment">false</bool>
  <bool name="blockingRefresh">true</bool>
  <bool name="useFromSideThreads">true</bool>
  <long name="sweepSamplingInterval">60</long>
</queryParser>

Segment-level Parallelism

The algorithm utilizes multiple threads on both sides if multiple threads are available. So, set indexSearcherExecutorThreads to -1 or >0 in the solr.xml file. Note the emphasis on both above: using the multiThreaded request parameter confines the inner (from-side) query to a single thread, limiting performance.

Sizing the From Side and the JVM

What the sidecar costs is driven by how large the from-side segments are, not by how many documents a join actually matches. For each (from-segment, to-segment) pair it builds, the parser holds several structures sized by the from segment’s maxDoc: the set of matching from documents, the map from each from document to its field value, and the batch written to the sidecar, which spans one document per from document however few of them the pair matched. On a from segment of several million documents each of those runs to tens of megabytes, so the two settings below are worth applying before a large join is put under sustained load.

Both are starting points drawn from benchmarking a single relation, not tuned defaults. Like the rest of this parser they are provisional; see Disclaimer.

Cap the From-Side Segment Size

Put a ceiling on how large the from core’s segments may grow, in that core’s solrconfig.xml:

<indexConfig>
  <mergePolicyFactory class="org.apache.solr.index.TieredMergePolicyFactory">
    <double name="maxMergedSegmentMB">200.0</double>
  </mergePolicyFactory>
</indexConfig>

This is a ceiling on future merges rather than a reduction: it prevents segments growing past the limit, but doesn’t split segments that are already larger.

TieredMergePolicy has no document-count setting, so this bounds documents only indirectly, through bytes. The document count a given maxMergedSegmentMB works out to depends on the average document size, and is worth re-checking whenever the schema changes.

Resist capping too aggressively. The parser maintains one pair column per (from-segment, to-segment) combination, so halving the from segment size doubles the number of from segments, and with it the number of pair columns to build, persist, and later reclaim.

Raise the G1 Region Size

GC_TUNE="-XX:G1HeapRegionSize=8m"

G1 treats any single allocation larger than half a region as “humongous”, placing it in its own run of regions that ordinary young collections don’t reclaim. At G1’s default 1 MB region that threshold is 512 KB, which the per-segment structures described above cross on segments of a few million documents — a bitset over a from segment passes it at roughly 4.2 million documents. An 8 MB region moves the threshold to 4 MB, which brings those allocations back under normal collection.

See JVM Settings for where JVM settings live and how to set them.

Disclaimer

auxIndexJoin is completely experimental, which is why it requires explicit configuration. So, far it’s evalueated within single relation (from-to index pair) only, but it should support many ones in principle. It’s parameters/usage is more subject to change with less backwards compatibility concern.