CASSANALYTICS-31 SAI index support in analytics - #220
Conversation
|
|
||
| // This builder always produces an index-less table (indexes are applied later by the 5.0 bridge), so a | ||
| // rebuild never carries indexes even if the caller passed index statements. buildSchema runs repeatedly per | ||
| // table in a JVM; if an earlier call already registered indexes, copy them onto this rebuild so it matches | ||
| // the registered table. | ||
| if (maybeExistingTableMetadata != null | ||
| && !maybeExistingTableMetadata.indexes.isEmpty() | ||
| && tableMetadata.indexes.isEmpty()) | ||
| { | ||
| tableMetadata = tableMetadata.unbuild() | ||
| .indexes(maybeExistingTableMetadata.indexes) | ||
| .build(); | ||
| } | ||
|
|
There was a problem hiding this comment.
RecordWriter is using 5 arg buildSchema() method which is not attaching index statements, you could eliminate copying indexes here if you use the 8 arg buildSchema() method in RecordWriter and can eventually get rid of 5 arg buildSchema() method
cqlTable = writerContext.bridge()
.buildSchema(writerContext.schema().getTableSchema().createStatement,
writerContext.job().qualifiedTableName().keyspace(),
IGNORED_REPLICATION_FACTOR,
writerContext.cluster().getPartitioner(),
writerContext.schema().getUserDefinedTypeStatements(),
null,
writerContext.schema().getTableSchema().getIndexStatements(),
false);
There was a problem hiding this comment.
I did restructure and moved SAI specific code to extended class. Now this comment not applicable I believe
There was a problem hiding this comment.
@skoppu22 This block of code is moved to FiveZeroSchemaBuilder.beforeTableRegistered but the issue still exists. You can get rid of the logic to not drop indexes in beforeTableRegistered and even the method beforeTableRegistered if you use 8 argument constructor for schemaBuilder in Recordwriter passing writerContext.schema().getTableSchema().getIndexStatements() in the constructor.
There was a problem hiding this comment.
@jyothsnakonisa I have applied your recommended changes to pass index statements as part of cqlTable. We still need beforeTableRegistered. Because
Reason 1, Some builders only receive a CREATE TABLE string, with no CqlTable and no CREATE INDEX statements, so they can only build index‑less. These are public CassandraBridge APIs:
encodePartitionKeys(...) / toTokens(...) / encodePartitionKey(...) → new FiveZeroSchemaBuilder(createTableStmt, ks, rf, partitioner)
readPartitionKeys(...) → same createStmt‑only builder
We can't make these pass statements without changing the public bridge API, and their callers frequently don't have the CREATE INDEX statements to pass. If one of these is the first build of a table in a JVM, the table registers without its SAI index.
Reason 2, The same table is built more than once per JVM in normal flows:
Write path, multiple partitions/executor: RecordWriter is constructed per Spark partition (mapPartitions); the 2nd+ partition on an executor rebuilds the already‑registered table.
Write path, single partition: after writing, SortedSSTableWriter.validateSSTables(...) → buildLocalDataLayer → new LocalDataLayer → buildSchema rebuilds the same table in the same JVM.
Read path: getCompactionScanner, getPartitionSizeIterator, rebuildBloomFilter all build from a CqlTable, repeatedly within a JVM.
There was a problem hiding this comment.
Looks like there are no callers for encodePartitionKeys and readPartitionKeys. If we can remove them, then we can switch to 8 arg constructor as you mentioned, and we can simplify beforeTableRegistered just to avoid duplicate registration
There was a problem hiding this comment.
Yes encodePartitionKeys & readPartitionKeys are not used anywhere in production code hence removing them should be safe. I will leave it to you whether to remove those methods or not.
There was a problem hiding this comment.
Marked them deprecated
| java.util.Collections.emptySet(), | ||
| java.util.Collections.emptySet()); |
There was a problem hiding this comment.
Please remove fully qualified class names.
|
|
||
| // This builder always produces an index-less table (indexes are applied later by the 5.0 bridge), so a | ||
| // rebuild never carries indexes even if the caller passed index statements. buildSchema runs repeatedly per | ||
| // table in a JVM; if an earlier call already registered indexes, copy them onto this rebuild so it matches | ||
| // the registered table. | ||
| if (maybeExistingTableMetadata != null | ||
| && !maybeExistingTableMetadata.indexes.isEmpty() | ||
| && tableMetadata.indexes.isEmpty()) | ||
| { | ||
| tableMetadata = tableMetadata.unbuild() | ||
| .indexes(maybeExistingTableMetadata.indexes) | ||
| .build(); | ||
| } | ||
|
|
There was a problem hiding this comment.
Yes encodePartitionKeys & readPartitionKeys are not used anywhere in production code hence removing them should be safe. I will leave it to you whether to remove those methods or not.
yifan-c
left a comment
There was a problem hiding this comment.
partial review; spotted a change that is needed. submitting the comments so far.
| if (componentFile.getFileName().toString().endsWith("Data.db")) | ||
| // send the primary data component ("<descriptor>-Data.db") last; SAI per-index components | ||
| // such as "...+TermsData.db" are still streamed here rather than skipped | ||
| if (SSTables.isDataComponent(componentFile)) |
There was a problem hiding this comment.
nit: Can you revert the comment?
This comment is already there in the javadoc of the method.
Repeating here and several other call-sites creates noise. The call-sites only need to know the method can determine whether a file is data component or not.
There was a problem hiding this comment.
This comment actually confused me a bit. It wasn't clear what "here" meant; I read it as "are still streamed here" then saw the immediate "if match: continue" logic which looks like the opposite of streaming here. :)
So if you drop the comment entirely, no worries. If you decide to keep it, I'd clarify something like "SAI per-index components such as "..." are still streamed below; we're only skipping core SSTable -Data.db files here."
There was a problem hiding this comment.
Reworded as suggested by Josh
| private static final Pattern COMPACTION_STRATEGY_PATTERN = Pattern.compile("compaction\\s*=\\s*\\{\\s*'class'\\s*:\\s*'([^']+)'"); | ||
|
|
||
| private static final Pattern MULTI_WHITESPACE_PATTERN = Pattern.compile("\\s+"); | ||
| private static final Pattern SAI_USING_PATTERN = Pattern.compile("USING '([^']*\\.)?STORAGEATTACHEDINDEX'"); |
There was a problem hiding this comment.
The pattern is not comprehensive. See the example below. cc: @maedhroz
cqlsh> DESC TABLE ks.tstsai;
CREATE TABLE ks.tstsai (
id int PRIMARY KEY,
val1 int,
val2 int,
val3 int,
val4 int,
val5 int
);
CREATE CUSTOM INDEX val1_index ON ks.tstsai (val1) USING 'STORAGEATTACHEDINDEX';
CREATE CUSTOM INDEX val2_index ON ks.tstsai (val2) USING 'storageattachedindex';
CREATE CUSTOM INDEX val3_index ON ks.tstsai (val3) USING 'sai';
CREATE CUSTOM INDEX val4_index ON ks.tstsai (val4) USING 'StorageAttachedIndex';
CREATE CUSTOM INDEX val5_index ON ks.tstsai (val5) USING 'org.apache.cassandra.index.sai.StorageAttachedIndex';
There was a problem hiding this comment.
Done, also added tests for all these cases
| public final class BroadcastableTableSchema implements Serializable | ||
| { | ||
| private static final long serialVersionUID = 1L; | ||
| private static final long serialVersionUID = 2L; |
There was a problem hiding this comment.
I understand the rational, but is it really needed? In the past we have made changes in class fields without updating serial version ID.
There was a problem hiding this comment.
The serialVersionUID change is unnecessary.
The serialization is not used in any long running service and no mixed-mode is expected.
jmckenzie-dev
left a comment
There was a problem hiding this comment.
Most of the way through review; will try and knock out the rest of it tomorrow.
| if (componentFile.getFileName().toString().endsWith("Data.db")) | ||
| // send the primary data component ("<descriptor>-Data.db") last; SAI per-index components | ||
| // such as "...+TermsData.db" are still streamed here rather than skipped | ||
| if (SSTables.isDataComponent(componentFile)) |
There was a problem hiding this comment.
This comment actually confused me a bit. It wasn't clear what "here" meant; I read it as "are still streamed here" then saw the immediate "if match: continue" logic which looks like the opposite of streaming here. :)
So if you drop the comment entirely, no worries. If you decide to keep it, I'd clarify something like "SAI per-index components such as "..." are still streamed below; we're only skipping core SSTable -Data.db files here."
# Conflicts: # CHANGES.txt # cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlTable.java # cassandra-bridge/src/testFixtures/java/org/apache/cassandra/spark/utils/test/TestSchema.java # cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
jmckenzie-dev
left a comment
There was a problem hiding this comment.
There are a lot of places where we call the old CassandraBridge#buildSchema where we can simplify to the signature assuming null UUID and false on enableCdc which would tidy up the patch a bit.
| public Set<String> indexStatements() | ||
| { | ||
| return indexCount; | ||
| return indexStatements; |
There was a problem hiding this comment.
If our goal is to keep external consumers from mutating the contents of this set we should do something like:
public Set<String> indexStatements() {
return Collections.unmodifiableSet(indexStatements);
}There was a problem hiding this comment.
indexStatements is already created using Collections.unmodifiableSet above in constructor
| * @param table the table name | ||
| * @return set of CREATE INDEX statements for the table | ||
| */ | ||
| public static Set<String> extractIndexStatements(@NotNull String schemaStr, |
There was a problem hiding this comment.
This caught my eye. I think there's a couple regex specific issues here; had a couple models take a couple passes at it and this is where we landed:
- Stray ? after the keyspace makes its last character optional. The format string ""?%s?"?" expands to "??"?, so the trailing ? binds the final character of the keyspace — asking for keyspace orders also matches order. Note the table side (correctly) has no ?; that asymmetry is the tell. This is copy-pasted from extractCleanedTableSchema, so it's pre-existing, but it's still wrong.
- Identifiers are interpolated raw into the pattern. A keyspace/table name containing a regex metacharacter (legal in quoted CQL identifiers) misbehaves or throws PatternSyntaxException. Wrap both in Pattern.quote(...).
- ⚠ These two are separate fixes. Pattern.quote alone does not fix bug CASSANDRA-18545: Provide a SecretsProvider interface to abstract the secret provisioning #1 — the ? lives in the format string outside the %s, so it survives quoting. We need to delete the ? and quote the value.
- Cosmetic: .{1} — the {1} is dead; it's just ..
Whenever I see big blocks of regex like this it always makes my "spidey senses tingle". Regexes are super powerful but are effectively embedded functions and edge-cases and "sanitization" of input become quite challenging when there's this much density expressed in one place.
| return quoteIdentifiers; | ||
| } | ||
|
|
||
| public Set<String> getIndexStatements() |
There was a problem hiding this comment.
Same here; returning a collection an external user could mutate under us. 😬
| + " weight float,\n" | ||
| + " height int\n" | ||
| + ");")); | ||
| + ");"), null, Collections.emptySet(), false); |
There was a problem hiding this comment.
Can update various call sites to simplify them now w/out the vestigial UUID and boolean now:
public CqlTable buildSchema(String createStatement,
String keyspace,
ReplicationFactor replicationFactor,
Partitioner partitioner,
Set<String> udts,
Set<String> indexStatements)| .withPartitioner(cassPartitioner) | ||
| .using(insertStatement) | ||
| // The data frame to write is always sorted, | ||
| // see org.apache.cassandra.spark.bulkwriter.CassandraBulkSourceRelation.insert |
There was a problem hiding this comment.
Is this a contractual promise from the other class or could this drift? There's no javadoc on CassandraBulkSourceRelation#insert and no javadoc on the parent class indicating that this is a durable "external API" commitment here.
There was a problem hiding this comment.
Added javadoc for why we need sorted
Generate SAI index files as well while producing sstable files, so no need of async rebuilding of SAI indexes after bulk write.
Circle CI : https://app.circleci.com/pipelines/github/skoppu22/cassandra-analytics/148/workflows/5eaeec43-e43b-444c-badb-ef696a3828fc