[
https://issues.apache.org/jira/browse/CASSANALYTICS-199?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18116199#comment-18116199
]
Caleb Rackliffe commented on CASSANALYTICS-199:
-----------------------------------------------
Let me know if you have any questions around the core SAI components you think
I might be able to help answer.
> Utilize SAI indexes for filtering in bulk read path
> ---------------------------------------------------
>
> Key: CASSANALYTICS-199
> URL: https://issues.apache.org/jira/browse/CASSANALYTICS-199
> Project: Apache Cassandra Analytics
> Issue Type: Improvement
> Components: Reader
> Reporter: Lukasz Antoniak
> Priority: Normal
>
> h2. Objective
> Improve bulk readers' performance by leveraging SAI indexes to scan only
> partition keys that potentially satisfy filter condition of the Spark data
> set. Currently, library supports only push-down filters for partition key
> columns. Filters based on non-partition key columns do not limit the result
> set returned from Cassandra Analytics to Spark. Appropriate filtering is
> performed by Spark engine.
> {code:java}
> bulkReaderDataFrame(TABLE).load()
> .filter("age > 30")
> .select("name")
> .collectAsList();{code}
> h2. Approach
> The bulk reader opens all SSTables participating in the requested token range
> and passes them through the normal compaction / reconciliation path.
> _SSTableReader_ already supports restricting the scan to a supplied list of
> partition key filters, which are applied consistently across the
> participating SSTables.
> The proposed SAI integration opens the relevant SAI indexes and creates a
> long-lived, lazy iterator of matching primary keys. Cassandra Analytics
> converts this stream into batches of unique partition keys, for example up to
> 1024 keys per batch. Each batch is then passed to {_}SSTableReader{_}, which
> reads only those candidate partitions from all participating SSTables before
> normal reconciliation is performed.
> If additional SAI matches remain, the same iterator is advanced to obtain the
> next batch, without reopening the SAI indexes.
> {code:java}
> Spark WHERE
> │
> ├── age > 30
> └── gender = 'female'
> │
> V
> QueryPlanner
> │
> ┌────────────┴───────────┐
> │ │
> age expression gender expression
> │ │
> V V
> For each SSTable:
> age segment 1 gender segment 1
> age segment 2 gender segment 2
> age segment 3 gender segment 3
> │ │
> V V
> UNION segments UNION segments
> │ │
> V V
> age result gender result
> per SSTable per SSTable
> │ │
> every predicate is evaluated
> independently across all SSTables
> │ │
> V V
> UNION across UNION across
> ALL SSTables ALL SSTables
> │ │
> └──────────┬────────────┘
> │
> V
> INTERSECT
> │
> │ lazy / long-lived
> │ KeyRangeIterator
> │
> V
> batch of up to 1024
> partition keys
> │
> V
> PartitionKeyFilters
> │
> V
> read candidate partitions from
> the participating Data.db SSTables
> │
> V
> CompactionStreamScanner
> reconciliation across SSTables
> tombstones / updates / TTLs
> │
> V
> reconstructed rows
> │
> V
> Spark filtering
> │
> V
> final result rows{code}
> h2. Implementation details
> # Unclear if we can re-use SAI {_}QueryController{_}, _QueryViewBuilder_ and
> {_}StorageAttachedIndexSearcher{_}, because those may query memtables and
> trigger read commands.
> # Re-use _IndexSegmentSearcher_ to read SAI index for given column and
> return partition keys. Opening multiple SAI segments needs to be reconciled
> with _KeyRangeUnionIterator_ before passing the keys to {_}SSTableReader{_}.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]