Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.txt
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
0.5.0
-----
* Mixed-case keyspace restore fails with keyspace does not exist when quoteIdentifiers is set (CASSANALYTICS-183)
* Opening one keyspace instance rebuilds every keyspace instance in the Cassandra 6.0 bridge (CASSANALYTICS-202)
* CDC logs NPE for deleted column values (CASSANALYTICS-178)
* Add Cassandra 6.0 support (CASSANALYTICS-37)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,14 @@ public String keyspace()
return keyspace;
}

/**
* @return the keyspace name, quoted with double quotes when {@code quoteIdentifiers} is set
*/
public String maybeQuotedKeyspace()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Wondering here if we really need wrappers with the "maybe*" prefix, or we can just use this as default for keyspace() and table(). What do you think @yifan-c / @frankgh ?

@bianca-stanciu29 bianca-stanciu29 Oct 2, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@bbotella @yifan-c @frankgh I think we should keep keyspace() and table() as raw-name accessors. They’re also used for SSTable metadata, local data-layer setup, and job statistics, so returning quoted names would require checking those callers too. The separate maybeQuoted accessors let us apply quoting when building Sidecar requests without changing other uses. CASSSIDECAR-475 fixes the metadata lookup inside Sidecar, but that’s separate from what these analytics accessors return. What do you think?

{
return maybeQuote(keyspace);
}

/**
* @return the table name in Cassandra
*/
Expand All @@ -73,6 +81,14 @@ public String table()
return table;
}

/**
* @return the table name, quoted with double quotes when {@code quoteIdentifiers} is set
*/
public String maybeQuotedTable()
{
return maybeQuote(table);
}

/**
* @return the identifiers should be quoted
*/
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.cassandra.spark.data;

import org.junit.jupiter.api.Test;

import static org.assertj.core.api.Assertions.assertThat;

class QualifiedTableNameTest
{
@Test
void testUnquotedIdentifiersReturnedAsIs()
{
QualifiedTableName name = new QualifiedTableName("mykeyspace", "mytable", false);
assertThat(name.keyspace()).isEqualTo("mykeyspace");
assertThat(name.table()).isEqualTo("mytable");
assertThat(name.maybeQuotedKeyspace()).isEqualTo("mykeyspace");
assertThat(name.maybeQuotedTable()).isEqualTo("mytable");
}

@Test
void testQuotedIdentifiersWrappedInDoubleQuotes()
{
QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true);
assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"MyKeyspace\"");
assertThat(name.maybeQuotedTable()).isEqualTo("\"MyTable\"");
}

@Test
void testRawAccessorsUnaffectedByQuoteFlag()
{
QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true);
assertThat(name.keyspace()).isEqualTo("MyKeyspace");
assertThat(name.table()).isEqualTo("MyTable");
}

@Test
void testToStringQuotesBothWhenFlagSet()
{
QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable", true);
assertThat(name.toString()).isEqualTo("\"MyKeyspace\".\"MyTable\"");
}

@Test
void testToStringUnquotedWhenFlagNotSet()
{
QualifiedTableName name = new QualifiedTableName("mykeyspace", "mytable", false);
assertThat(name.toString()).isEqualTo("mykeyspace.mytable");
}

@Test
void testDefaultConstructorDoesNotQuote()
{
QualifiedTableName name = new QualifiedTableName("MyKeyspace", "MyTable");
assertThat(name.quoteIdentifiers()).isFalse();
assertThat(name.maybeQuotedKeyspace()).isEqualTo("MyKeyspace");
assertThat(name.maybeQuotedTable()).isEqualTo("MyTable");
}

@Test
void testReservedWordIdentifiersQuoted()
{
QualifiedTableName name = new QualifiedTableName("keyspace", "table", true);
assertThat(name.maybeQuotedKeyspace()).isEqualTo("\"keyspace\"");
assertThat(name.maybeQuotedTable()).isEqualTo("\"table\"");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -289,11 +289,12 @@ protected String getCurrentKeyspaceSchema() throws Exception
private TokenRangeReplicasResponse getTokenRangesAndReplicaSets()
{
CassandraContext context = getCassandraContext();
String maybeQuotedKeyspace = maybeQuotedIdentifier(bridge(), conf.quoteIdentifiers, conf.keyspace);
try
{
long start = System.nanoTime();
TokenRangeReplicasResponse response = context.getSidecarClient()
.tokenRangeReplicas(new ArrayList<>(context.getCluster()), conf.keyspace)
.tokenRangeReplicas(new ArrayList<>(context.getCluster()), maybeQuotedKeyspace)
.get();
long elapsedTimeNanos = System.nanoTime() - start;
LOGGER.info("Retrieved token ranges for {} instances in {} milliseconds",
Expand All @@ -303,8 +304,8 @@ private TokenRangeReplicasResponse getTokenRangesAndReplicaSets()
}
catch (ExecutionException | InterruptedException exception)
{
LOGGER.error("Failed to get token ranges for keyspace {}", conf.keyspace, exception);
throw new SidecarApiCallException("Failed to get token ranges for keyspace" + conf.keyspace, exception);
LOGGER.error("Failed to get token ranges for keyspace {}", maybeQuotedKeyspace, exception);
throw new SidecarApiCallException("Failed to get token ranges for keyspace " + maybeQuotedKeyspace, exception);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,8 +103,8 @@ public void createRestoreJob(CreateRestoreJobRequestPayload createRestoreJobRequ
try
{
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
sidecarClient.createRestoreJob(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
sidecarClient.createRestoreJob(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
createRestoreJobRequestPayload).get();
}
catch (Exception exception)
Expand All @@ -121,8 +121,8 @@ public RestoreJobSummaryResponsePayload restoreJobSummary() throws SidecarApiCal
try
{
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
return sidecarClient.restoreJobSummary(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
return sidecarClient.restoreJobSummary(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
jobId).get();
}
catch (Exception exception)
Expand Down Expand Up @@ -168,8 +168,8 @@ public void updateRestoreJob(UpdateRestoreJobRequestPayload updateRestoreJobRequ
{
LOGGER.info("Updating the restore job. clusterId={} restoreJobId={}", clusterId, jobId);
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
sidecarClient.updateRestoreJob(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
sidecarClient.updateRestoreJob(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
jobId,
updateRestoreJobRequestPayload).get();
}
Expand All @@ -188,8 +188,8 @@ public void abortRestoreJob() throws SidecarApiCallException
{
LOGGER.info("Abort job. clusterId={} restoreJobId={}", clusterId, jobId);
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
sidecarClient.abortRestoreJob(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
sidecarClient.abortRestoreJob(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
jobId).get();
}
catch (Exception exception)
Expand All @@ -207,8 +207,8 @@ private CompletableFuture<Void> createRestoreSliceWithCustomRetry(SidecarInstanc
RetryPolicy retryPolicy)
{
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
jobInfo.getRestoreJobId(clusterId),
createSliceRequestPayload);
return sidecarClient.executeRequestAsync(sidecarClient.requestBuilder()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -169,8 +169,8 @@ private void createRestoreSliceInternal(String clusterId,
SidecarClient sidecarClient = dataTransferApi.sidecarClient();
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
UUID restoreJobId = jobInfo.getRestoreJobId(clusterId);
CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
CreateRestoreJobSliceRequest request = new CreateRestoreJobSliceRequest(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
restoreJobId,
createSliceRequestPayload);
RetryPolicy retryPolicy = new CloudStorageDataTransferApiImpl.ExecutorCreateSliceRetryPolicy(sidecarClient);
Expand All @@ -197,8 +197,8 @@ private void restoreJobProgressInternal(String clusterId,
JobInfo jobInfo = dataTransferApi.jobInfo();
QualifiedTableName qualifiedTableName = jobInfo.qualifiedTableName();
UUID restoreJobId = jobInfo.getRestoreJobId(clusterId);
RestoreJobProgressRequestParams requestParams = new RestoreJobProgressRequestParams(qualifiedTableName.keyspace(),
qualifiedTableName.table(),
RestoreJobProgressRequestParams requestParams = new RestoreJobProgressRequestParams(qualifiedTableName.maybeQuotedKeyspace(),
qualifiedTableName.maybeQuotedTable(),
restoreJobId,
fetchPolicy);
RestoreJobProgressResponsePayload jobProgress = restoreJobProgress(dataTransferApi, requestParams);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1189,7 +1189,9 @@ protected Sizing getSizing(CompletableFuture<RingResponse> ringFuture,
ReplicationFactor replicationFactor,
ClientConfig options)
{
return SizingFactory.create(replicationFactor, options, consistencyLevel, keyspace, table, datacenter, sidecar, sidecarPort, ringFuture);
return SizingFactory.create(replicationFactor, options, consistencyLevel,
maybeQuotedKeyspace, maybeQuotedTable, datacenter,
sidecar, sidecarPort, ringFuture);
}

protected void await(CountDownLatch latch)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,33 +24,87 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.TreeMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;

import com.google.common.collect.ImmutableMap;
import com.google.common.util.concurrent.Uninterruptibles;
import org.apache.spark.SparkConf;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;

import o.a.c.sidecar.client.shaded.common.response.NodeSettings;
import o.a.c.sidecar.client.shaded.common.response.SchemaResponse;
import o.a.c.sidecar.client.shaded.common.response.TimeSkewResponse;
import o.a.c.sidecar.client.shaded.common.response.TokenRangeReplicasResponse;
import o.a.c.sidecar.client.shaded.client.SidecarClient;
import org.apache.cassandra.bridge.CassandraBridge;
import org.apache.cassandra.spark.bulkwriter.token.TokenRangeMapping;
import org.apache.cassandra.spark.exception.TimeSkewTooLargeException;

import static org.apache.cassandra.spark.TestUtils.range;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

public class CassandraClusterInfoTest
{
@Test
void testQuotedKeyspacePassedToSchemaAndTokenRanges() throws Exception
{
Map<String, String> options = new TreeMap<>(String.CASE_INSENSITIVE_ORDER);
options.put(WriterOptions.SIDECAR_INSTANCES.name(), "127.0.0.1");
options.put(WriterOptions.KEYSPACE.name(), "MyKeyspace");
options.put(WriterOptions.TABLE.name(), "MyTable");
options.put(WriterOptions.QUOTE_IDENTIFIERS.name(), "true");
BulkSparkConf conf = new BulkSparkConf(new SparkConf(), options);
CassandraBridge bridge = mock(CassandraBridge.class);
when(bridge.maybeQuoteIdentifier("MyKeyspace")).thenReturn("\"MyKeyspace\"");
CassandraContext context = mock(CassandraContext.class);
SidecarClient sidecar = mock(SidecarClient.class);
when(context.getSidecarClient()).thenReturn(sidecar);
when(context.getCluster()).thenReturn(Collections.emptySet());
SchemaResponse schema = mock(SchemaResponse.class);
when(schema.schema()).thenReturn("schema");
when(sidecar.schema("\"MyKeyspace\"")).thenReturn(CompletableFuture.completedFuture(schema));
CompletableFuture<TokenRangeReplicasResponse> failed = new CompletableFuture<>();
failed.completeExceptionally(new IllegalStateException("test failure"));
when(sidecar.tokenRangeReplicas(any(), eq("\"MyKeyspace\""))).thenReturn(failed);

try (CassandraClusterInfo ci = new CassandraClusterInfo(conf)
{
@Override
protected CassandraContext buildCassandraContext()
{
return context;
}

@Override
protected CassandraBridge bridge()
{
return bridge;
}
})
{
assertThat(ci.getCurrentKeyspaceSchema()).isEqualTo("schema");
assertThatThrownBy(() -> ci.getTokenRangeMapping(false)).hasMessageContaining("Unable to initialize ring information");
verify(sidecar).schema("\"MyKeyspace\"");
verify(sidecar).tokenRangeReplicas(any(), eq("\"MyKeyspace\""));
}
}

@Test
void testTimeSkewAcceptable()
{
Expand Down
Loading