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
2 changes: 1 addition & 1 deletion spark/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ Prerequisites

* Java 8 (the module targets source/target 1.8).
* Maven 3.8+.
* Access to a Spark 3.4.x cluster with the Skyflow Java SDK (3.0.0-beta.6) compatible dependencies.
* Access to a Spark 3.4.x cluster with the Skyflow Java SDK (skyflow-flowvault-java, 1.1.0) compatible dependencies.
* Skyflow vault credentials stored as a JSON string or resolvable via your secret manager.

Building The Jar
Expand Down
2 changes: 1 addition & 1 deletion spark/dependency-reduced-pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@
</filters>
<artifactSet>
<includes>
<include>com.skyflow:skyflow-java</include>
<include>com.skyflow:skyflow-flowvault-java</include>
</includes>
</artifactSet>
<createDependencyReducedPom>true</createDependencyReducedPom>
Expand Down
6 changes: 3 additions & 3 deletions spark/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,7 @@
</filters>
<artifactSet>
<includes>
<include>com.skyflow:skyflow-java</include>
<include>com.skyflow:skyflow-flowvault-java</include>
</includes>
</artifactSet>
<createDependencyReducedPom>true</createDependencyReducedPom>
Expand Down Expand Up @@ -125,8 +125,8 @@
</dependency>
<dependency>
<groupId>com.skyflow</groupId>
<artifactId>skyflow-java</artifactId>
<version>3.0.0-beta.8</version>
<artifactId>skyflow-flowvault-java</artifactId>
<version>1.1.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
Expand Down
1 change: 1 addition & 0 deletions spark/src/main/java/com/skyflow/spark/Constants.java
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ public class Constants {
public static final long MAX_DELAY_MILLI_SECONDS = 10000; // max delay
// Set of retryable error codes
public static final HashSet<Integer> RETRYABLE_ERROR_CODES = new HashSet<>(Arrays.asList(
409,
429,
500,
502,
Expand Down
156 changes: 86 additions & 70 deletions spark/src/main/java/com/skyflow/spark/Helper.java

Large diffs are not rendered by default.

61 changes: 39 additions & 22 deletions spark/src/main/java/com/skyflow/spark/VaultHelper.java
Original file line number Diff line number Diff line change
Expand Up @@ -5,12 +5,13 @@
import com.skyflow.config.Credentials;
import com.skyflow.errors.SkyflowException;
import com.skyflow.vault.data.ErrorRecord;
import com.skyflow.vault.data.InsertResponse;
import com.skyflow.vault.data.InsertRequest;
import com.skyflow.vault.data.DetokenizeRequest;
import com.skyflow.vault.data.DetokenizeResponse;
import com.skyflow.vault.data.DetokenizeResponseObject;
import com.skyflow.vault.data.Success;
import com.skyflow.vault.data.BulkInsertRequest;
import com.skyflow.vault.data.BulkInsertResponse;
import com.skyflow.vault.data.BulkInsertResponseRecord;
import com.skyflow.vault.data.InsertRequestRecord;
import com.skyflow.vault.data.BulkDetokenizeRequest;
import com.skyflow.vault.data.BulkDetokenizeResponse;
import com.skyflow.vault.data.BulkDetokenizeResponseRecord;

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
Expand Down Expand Up @@ -73,7 +74,7 @@ public void initializeSkyflowClient(TableHelper tableHelper) throws SkyflowExcep
VaultConfig vaultConfig = new VaultConfig();
vaultConfig.setVaultId(tableHelper.getVaultId());
vaultConfig.setClusterId(tableHelper.getClusterId());
vaultConfig.setVaultURL(tableHelper.getVaultUrl());
vaultConfig.setVaultUrl(tableHelper.getVaultUrl());
vaultConfig.setEnv(tableHelper.getEnv());
vaultConfig.setCredentials(credentials);
this.skyflowClient = getSkyflowBuilder()
Expand Down Expand Up @@ -129,16 +130,16 @@ public Dataset<Row> tokenize(TableHelper tableHelper, Dataset<Row> dataToIngest,
for (List<Row> batch : Helper.getBatches(dataToIngest, batchSize)) {
List<Row> batchOutputRows;
// Construct and send insert request
InsertRequest insertRequest = Helper.constructInsertRequest(schemaMappings, batch);
BulkInsertRequest insertRequest = Helper.constructInsertRequest(schemaMappings, batch);
if(insertRequest.getRecords().isEmpty()) {
batchOutputRows = Helper.replaceDataWithTokens(schemaMappings, batch, new HashMap<>(), new HashMap<>());
} else {
logger.info(LOG_PREFIX + "Processing batch #" + batchNumber + ", No.of records: "
+ insertRequest.getRecords().size());
InsertResponse insertResponse = skyflowClient.vault().bulkInsert(insertRequest);
BulkInsertResponse insertResponse = skyflowClient.vault().bulkInsert(insertRequest);

// Process success and error responses
Map<Object, Success> successMap = Helper.getInsertSuccessMap(insertResponse,
Map<Object, BulkInsertResponseRecord> successMap = Helper.getInsertSuccessMap(insertResponse,
insertRequest.getRecords());
Map<Object, ErrorRecord> errorsMap = Helper.getInsertErrorsMap(insertResponse,
insertRequest.getRecords());
Expand Down Expand Up @@ -199,19 +200,24 @@ public Dataset<Row> detokenize(TableHelper tableHelper, Dataset<Row> tokenizedDa
for (List<Row> batch : Helper.getBatches(tokenizedData, batchSize)) {
List<Row> batchOutputRows;
// Construct and send detokenize request
DetokenizeRequest detokenizeRequest = Helper.constructDetokenizeRequest(schemaMappings, batch);
BulkDetokenizeRequest detokenizeRequest = Helper.constructDetokenizeRequest(schemaMappings, batch);
if(detokenizeRequest.getTokens().isEmpty()) {
batchOutputRows = Helper.replaceTokensWithData(schemaMappings, batch, new HashMap<>(), new HashMap<>());
} else {
logger.info(LOG_PREFIX + "Processing batch #" + batchNumber + ", No.of records: "
+ detokenizeRequest.getTokens().size());
DetokenizeResponse detokenizeResponse = skyflowClient.vault().bulkDetokenize(detokenizeRequest);
BulkDetokenizeResponse detokenizeResponse = skyflowClient.vault().bulkDetokenize(detokenizeRequest);

// Process success and error responses
Map<String, DetokenizeResponseObject> successMap = Helper.getDetokenizeSuccessMap(detokenizeResponse);
Map<String, ErrorRecord> errorsMap = Helper.geDetokenizeErrorsMap(detokenizeResponse,
Map<String, BulkDetokenizeResponseRecord> successMap = Helper.getDetokenizeSuccessMap(detokenizeResponse);
Map<String, ErrorRecord> errorsMap = Helper.getDetokenizeErrorsMap(detokenizeResponse,
detokenizeRequest.getTokens());
logger.fine(LOG_PREFIX + "Success count: " + successMap.size() + " Error count: " + errorsMap.size());
if (successMap.size() + errorsMap.size() != detokenizeRequest.getTokens().size()) {
logger.warning(LOG_PREFIX + "Detokenize response accounted for "
+ (successMap.size() + errorsMap.size()) + " of " + detokenizeRequest.getTokens().size()
+ " requested tokens; some tokens got no response entry.");
}
// Retry failed tokens if necessary
if (detokenizeResponse.getSummary().getTotalFailed() > 0) {
retryFailedTokens(detokenizeRequest, successMap, errorsMap);
Expand All @@ -232,22 +238,26 @@ public Dataset<Row> detokenize(TableHelper tableHelper, Dataset<Row> tokenizedDa
}

// Helper method to retry failed records with exponential backoff and jitter
private void retryFailedRecords(InsertRequest request,
Map<Object, Success> successMap,
private void retryFailedRecords(BulkInsertRequest request,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can we use BulkInsertResponse.getRecordsToRetry() here

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

public static final HashSet<Integer> RETRYABLE_ERROR_CODES = new HashSet<>(Arrays.asList(

Does it return records that cover these error codes @Devesh-Skyflow ?
I believe 429 is not covered.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Yes 429 is not covered

Map<Object, BulkInsertResponseRecord> successMap,
Map<Object, ErrorRecord> errorsMap) throws SkyflowException {
int currentRetry = 0;
// errorsMap indices are relative to whichever batch was most recently sent, not the
// original request, so the lookup list must track the current batch across rounds.
List<InsertRequestRecord> currentBatch = request.getRecords();
while (!errorsMap.isEmpty() && currentRetry < retryCount) {
InsertRequest retryRequest = Helper.constructInsertRetryRequest(request.getRecords(), errorsMap);
BulkInsertRequest retryRequest = Helper.constructInsertRetryRequest(currentBatch, errorsMap);
if (retryRequest.getRecords().isEmpty()) {
logger.fine(LOG_PREFIX + NO_RETRIES_NEEDED_PROCEEDING);
break;
} else {
logger.fine(
LOG_PREFIX + "Retrying " + retryRequest.getRecords().size() + " failed records. Attempt: "
+ (currentRetry + 1));
InsertResponse retryResponse = skyflowClient.vault().bulkInsert(retryRequest);
BulkInsertResponse retryResponse = skyflowClient.vault().bulkInsert(retryRequest);
Helper.sleepWithExponentialBackoff(currentRetry);
Helper.mergeInsertRetryResults(retryRequest.getRecords(), retryResponse, successMap, errorsMap);
currentBatch = retryRequest.getRecords();
logger.fine(LOG_PREFIX + "After retry, Success count: " + successMap.size() + " Error count: "
+ errorsMap.size());
currentRetry++;
Expand All @@ -256,8 +266,8 @@ private void retryFailedRecords(InsertRequest request,
}

// Helper method to retry failed tokens with exponential backoff and jitter
private void retryFailedTokens(DetokenizeRequest detokenizeRequest,
Map<String, DetokenizeResponseObject> successMap, Map<String, ErrorRecord> errorsMap)
private void retryFailedTokens(BulkDetokenizeRequest detokenizeRequest,
Map<String, BulkDetokenizeResponseRecord> successMap, Map<String, ErrorRecord> errorsMap)
throws SkyflowException {
int currentRetry = 0;
while (!errorsMap.isEmpty() && currentRetry < retryCount) {
Expand All @@ -273,12 +283,19 @@ private void retryFailedTokens(DetokenizeRequest detokenizeRequest,
LOG_PREFIX + "Retrying " + retryableTokens.size() + " failed tokens. Attempt: "
+ (currentRetry + 1));
Helper.sleepWithExponentialBackoff(currentRetry);
DetokenizeResponse retryResponse = skyflowClient.vault()
.bulkDetokenize(DetokenizeRequest.builder().tokens(retryableTokens)
BulkDetokenizeResponse retryResponse = skyflowClient.vault()
.bulkDetokenize(BulkDetokenizeRequest.builder().tokens(retryableTokens)
.tokenGroupRedactions(detokenizeRequest.getTokenGroupRedactions()).build());
Helper.mergeDetokenizeRetryResults(retryResponse, retryableTokens, successMap, errorsMap);
logger.fine(LOG_PREFIX + "After retry, Success count: " + successMap.size() + " Error count: "
+ errorsMap.size());
long unaccountedForTokens = retryableTokens.stream()
.filter(token -> !successMap.containsKey(token) && !errorsMap.containsKey(token))
.count();
if (unaccountedForTokens > 0) {
logger.warning(LOG_PREFIX + unaccountedForTokens + " of " + retryableTokens.size()
+ " retried tokens got no response entry (success or error) after this retry.");
}
currentRetry++;
}
}
Expand Down
Loading