[ 
https://issues.apache.org/jira/browse/CASSANALYTICS-199?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18116286#comment-18116286
 ] 

Piotr Kolaczkowski commented on CASSANALYTICS-199:
--------------------------------------------------

{quote}Cassandra Analytics converts this stream into batches of unique 
partition keys, for example up to 1024 keys per batch. 
{quote}
Not sure, why do you need batching of fixed size? Getting 1024 keys from SAI 
does not necessarily mean a sequential scan underneath. It can fetch many more 
if intersections are involved.
You may want to batch as many keys as needed to be efficient on the SSTable 
reading side. So maybe collecting all keys for the whole partition? 


{quote}Unclear if we can re-use SAI {_}QueryController{_}, _QueryViewBuilder_ 
and {_}StorageAttachedIndexSearcher{_}, because those may query memtables and 
trigger read commands.
{quote}
You likely want to write your own `ResultRetriever` and use the rest of SAI 
as-is. See the existing ResultRetriever(s).


{quote}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{_}.
{quote}
That looks very much like implementing 80% of SAI by yourself ;)I don't think 
you want to go that route. 

> 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]

Reply via email to