diff --git a/spark/README.md b/spark/README.md
index 87cfde6..573b0e8 100644
--- a/spark/README.md
+++ b/spark/README.md
@@ -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
diff --git a/spark/dependency-reduced-pom.xml b/spark/dependency-reduced-pom.xml
index 4895b78..71f9a79 100644
--- a/spark/dependency-reduced-pom.xml
+++ b/spark/dependency-reduced-pom.xml
@@ -74,7 +74,7 @@
- com.skyflow:skyflow-java
+ com.skyflow:skyflow-flowvault-java
true
diff --git a/spark/pom.xml b/spark/pom.xml
index ac505b1..647d98b 100644
--- a/spark/pom.xml
+++ b/spark/pom.xml
@@ -88,7 +88,7 @@
- com.skyflow:skyflow-java
+ com.skyflow:skyflow-flowvault-java
true
@@ -125,8 +125,8 @@
com.skyflow
- skyflow-java
- 3.0.0-beta.8
+ skyflow-flowvault-java
+ 1.1.0
org.apache.spark
diff --git a/spark/src/main/java/com/skyflow/spark/Constants.java b/spark/src/main/java/com/skyflow/spark/Constants.java
index 351fcd6..cbe1db4 100644
--- a/spark/src/main/java/com/skyflow/spark/Constants.java
+++ b/spark/src/main/java/com/skyflow/spark/Constants.java
@@ -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 RETRYABLE_ERROR_CODES = new HashSet<>(Arrays.asList(
+ 409,
429,
500,
502,
diff --git a/spark/src/main/java/com/skyflow/spark/Helper.java b/spark/src/main/java/com/skyflow/spark/Helper.java
index 568c542..570be6a 100644
--- a/spark/src/main/java/com/skyflow/spark/Helper.java
+++ b/spark/src/main/java/com/skyflow/spark/Helper.java
@@ -3,14 +3,17 @@
import com.fasterxml.jackson.core.type.TypeReference;
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.InsertRecord;
-import com.skyflow.vault.data.DetokenizeRequest;
-import com.skyflow.vault.data.DetokenizeResponse;
-import com.skyflow.vault.data.DetokenizeResponseObject;
+import com.skyflow.vault.data.BulkInsertRequest;
+import com.skyflow.vault.data.BulkInsertRequestRecord;
+import com.skyflow.vault.data.BulkInsertResponse;
+import com.skyflow.vault.data.BulkInsertResponseRecord;
+import com.skyflow.vault.data.InsertRequestRecord;
+import com.skyflow.vault.data.InsertResponseRecord;
+import com.skyflow.vault.data.BulkDetokenizeRequest;
+import com.skyflow.vault.data.BulkDetokenizeResponse;
+import com.skyflow.vault.data.BulkDetokenizeResponseRecord;
import com.skyflow.vault.data.TokenGroupRedactions;
-import com.skyflow.vault.data.Success;
+import com.skyflow.vault.data.UpsertOptions;
import com.skyflow.vault.data.Token;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
@@ -147,24 +150,24 @@ public List next() {
* Tokenize util methods
*/
- // Constructs an InsertRequest object from a batch of rows and column mappings
- public static InsertRequest constructInsertRequest(Map schemaMappings, List batch) {
- ArrayList records = new ArrayList<>();
+ // Constructs a BulkInsertRequest object from a batch of rows and column mappings
+ public static BulkInsertRequest constructInsertRequest(Map schemaMappings, List batch) {
+ ArrayList records = new ArrayList<>();
// Track seen values per table + vault column
Map>> valuesDedupMap = new HashMap<>();
for (Row row : batch) {
- List rowRecords = constructInsertRecordsForRow(row, valuesDedupMap, schemaMappings);
+ List rowRecords = constructInsertRecordsForRow(row, valuesDedupMap, schemaMappings);
records.addAll(rowRecords);
}
- return InsertRequest.builder()
+ return BulkInsertRequest.builder()
.records(records)
.build();
}
- private static List constructInsertRecordsForRow(Row row,
+ private static List constructInsertRecordsForRow(Row row,
Map>> seenValues, Map schemaMappings) {
- List records = new ArrayList<>();
+ List records = new ArrayList<>();
for (Map.Entry entry : schemaMappings.entrySet()) {
String datasetColumn = entry.getKey();
@@ -192,48 +195,57 @@ private static List constructInsertRecordsForRow(Row row,
HashMap record = new HashMap<>();
record.put(vaultColumn, row.getAs(datasetColumn));
if(skyflowColumnMapping.getIsUnique() != null && skyflowColumnMapping.getIsUnique() == true) {
- records.add(InsertRecord.builder().data(record).table(skyflowColumnMapping.getTableName())
- .upsert(Collections.singletonList(skyflowColumnMapping.getColumnName())).build());
+ records.add(BulkInsertRequestRecord.builder().data(record).tableName(skyflowColumnMapping.getTableName())
+ .upsert(UpsertOptions.builder()
+ .uniqueColumns(Collections.singletonList(skyflowColumnMapping.getColumnName()))
+ .build())
+ .build());
} else {
- records.add(InsertRecord.builder().data(record).table(skyflowColumnMapping.getTableName()).build());
+ records.add(BulkInsertRequestRecord.builder().data(record).tableName(skyflowColumnMapping.getTableName()).build());
}
}
return records;
}
- // Converts InsertResponse success records into a map for quick lookup
- public static Map