diff --git a/build.gradle b/build.gradle
index 7c72f589..62c94c98 100644
--- a/build.gradle
+++ b/build.gradle
@@ -15,14 +15,16 @@
*/
buildscript {
ext {
- dataflowBeamVersion = '2.23.0'
+ dataflowBeamVersion = '2.26.0'
}
repositories {
mavenCentral()
jcenter()
maven {
url "https://plugins.gradle.org/m2/"
+
}
+
dependencies {
classpath "net.ltgt.gradle:gradle-apt-plugin:0.19"
classpath "com.diffplug.spotless:spotless-plugin-gradle:3.24.2"
@@ -80,7 +82,9 @@ tasks.withType(JavaCompile) {
}
repositories {
- mavenCentral()
+
+ mavenCentral()
+
}
dependencies {
@@ -89,6 +93,7 @@ dependencies {
compile group: 'org.apache.beam', name: 'beam-runners-direct-java', version: dataflowBeamVersion
compile group: 'org.apache.beam', name: 'beam-sdks-java-extensions-ml', version: dataflowBeamVersion
compile group: 'org.apache.beam', name: 'beam-sdks-java-io-amazon-web-services', version: dataflowBeamVersion
+ compile group: 'org.apache.beam', name: 'beam-sdks-java-io-contextualtextio', version: dataflowBeamVersion
compile group: 'org.slf4j', name: 'slf4j-jdk14', version: '1.7.5'
compile 'com.google.cloud:google-cloud-kms:0.70.0-beta'
compile 'com.google.guava:guava:27.0-jre'
diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties
index 6b3851a8..a4b44297 100644
--- a/gradle/wrapper/gradle-wrapper.properties
+++ b/gradle/wrapper/gradle-wrapper.properties
@@ -1,5 +1,5 @@
distributionBase=GRADLE_USER_HOME
distributionPath=wrapper/dists
-distributionUrl=https\://services.gradle.org/distributions/gradle-5.1-bin.zip
+distributionUrl=https\://services.gradle.org/distributions/gradle-6.3-bin.zip
zipStoreBase=GRADLE_USER_HOME
zipStorePath=wrapper/dists
diff --git a/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2.java b/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2.java
index c8ea9915..e94f92dd 100644
--- a/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2.java
+++ b/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2.java
@@ -24,7 +24,6 @@
import com.google.swarm.tokenization.common.BigQueryDynamicWriteTransform;
import com.google.swarm.tokenization.common.BigQueryReadTransform;
import com.google.swarm.tokenization.common.BigQueryTableHeaderDoFn;
-import com.google.swarm.tokenization.common.CSVFileReaderSplitDoFn;
import com.google.swarm.tokenization.common.DLPTransform;
import com.google.swarm.tokenization.common.ExtractColumnNamesTransform;
import com.google.swarm.tokenization.common.FilePollingTransform;
@@ -43,11 +42,13 @@
import org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.io.FileIO;
import org.apache.beam.sdk.io.FileIO.ReadableFile;
+import org.apache.beam.sdk.io.contextualtextio.ContextualTextIO;
import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.Flatten;
import org.apache.beam.sdk.transforms.GroupByKey;
import org.apache.beam.sdk.transforms.ParDo;
+import org.apache.beam.sdk.transforms.Values;
import org.apache.beam.sdk.transforms.View;
import org.apache.beam.sdk.transforms.windowing.AfterProcessingTime;
import org.apache.beam.sdk.transforms.windowing.FixedWindows;
@@ -84,7 +85,7 @@ public static void main(String[] args) {
public static PipelineResult run(DLPTextToBigQueryStreamingV2PipelineOptions options) {
Pipeline p = Pipeline.create(options);
-
+ // p.apply(ContextualTextIO.readFiles());
switch (options.getDLPMethod()) {
case INSPECT:
case DEID:
@@ -123,19 +124,17 @@ public static PipelineResult run(DLPTextToBigQueryStreamingV2PipelineOptions opt
case CSV:
records =
inputFiles
+ .apply(Values.create())
.apply(
- "SplitCSVFile",
- ParDo.of(
- new CSVFileReaderSplitDoFn(
- options.getKeyRange(),
- options.getRecordDelimiter(),
- options.getSplitSize())))
+ ContextualTextIO.readFiles()
+ .withDelimiter(options.getRecordDelimiter().getBytes()))
.apply(
"ConvertToDLPRow",
ParDo.of(new ConvertCSVRecordToDLPRow(options.getColumnDelimiter(), header))
.withSideInputs(header));
break;
case JSON:
+
records =
inputFiles
.apply(
diff --git a/src/main/java/com/google/swarm/tokenization/beam/ConvertCSVRecordToDLPRow.java b/src/main/java/com/google/swarm/tokenization/beam/ConvertCSVRecordToDLPRow.java
index b47b5ca4..38664a43 100644
--- a/src/main/java/com/google/swarm/tokenization/beam/ConvertCSVRecordToDLPRow.java
+++ b/src/main/java/com/google/swarm/tokenization/beam/ConvertCSVRecordToDLPRow.java
@@ -20,10 +20,11 @@
import com.google.swarm.tokenization.common.Util;
import java.io.IOException;
import java.util.List;
-import java.util.Objects;
+import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollectionView;
+import org.apache.beam.sdk.values.Row;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -34,18 +35,13 @@
*
If a column delimiter of values isn't provided, input is assumed to be unstructured and the
* input KV value is saved in a single column of output {@link Table.Row}.
*/
-public class ConvertCSVRecordToDLPRow extends DoFn, KV> {
+public class ConvertCSVRecordToDLPRow extends DoFn> {
public static final Logger LOG = LoggerFactory.getLogger(ConvertCSVRecordToDLPRow.class);
private final Character columnDelimiter;
private PCollectionView> header;
- public ConvertCSVRecordToDLPRow(PCollectionView> header) {
- this.columnDelimiter = null;
- this.header = header;
- }
-
public ConvertCSVRecordToDLPRow(Character columnDelimiter, PCollectionView> header) {
this.columnDelimiter = columnDelimiter;
this.header = header;
@@ -54,16 +50,19 @@ public ConvertCSVRecordToDLPRow(Character columnDelimiter, PCollectionView csvHeader = context.sideInput(header);
if (columnDelimiter != null) {
-
List values = Util.parseLine(input, columnDelimiter, '"');
if (values.size() == csvHeader.size()) {
values.forEach(
value -> rowBuilder.addValues(Value.newBuilder().setStringValue(value).build()));
- context.output(KV.of(context.element().getKey(), rowBuilder.build()));
+ context.output(KV.of(filename, rowBuilder.build()));
} else {
LOG.warn(
@@ -73,7 +72,7 @@ public void processElement(ProcessContext context) throws IOException {
}
} else {
rowBuilder.addValues(Value.newBuilder().setStringValue(input).build());
- context.output(KV.of(context.element().getKey(), rowBuilder.build()));
+ context.output(KV.of(filename, rowBuilder.build()));
}
}
}
diff --git a/src/main/java/com/google/swarm/tokenization/common/BigQueryDynamicWriteTransform.java b/src/main/java/com/google/swarm/tokenization/common/BigQueryDynamicWriteTransform.java
index d3ebfa8e..b44b44f4 100644
--- a/src/main/java/com/google/swarm/tokenization/common/BigQueryDynamicWriteTransform.java
+++ b/src/main/java/com/google/swarm/tokenization/common/BigQueryDynamicWriteTransform.java
@@ -90,7 +90,6 @@ public BQDestination(String datasetName, String projectId) {
public TableDestination getTable(KV destination) {
TableDestination dest =
new TableDestination(destination.getKey(), "DLP Transformation Storage Table");
- LOG.debug("Table Destination {}", dest.getTableSpec());
return dest;
}
@@ -98,14 +97,13 @@ public TableDestination getTable(KV destination) {
public KV getDestination(ValueInSingleWindow> element) {
String key = element.getValue().getKey();
String tableName = String.format("%s:%s.%s", projectId, datasetName, key);
- LOG.debug("Table Name {}", tableName);
return KV.of(tableName, element.getValue().getValue());
}
@Override
public TableSchema getSchema(KV destination) {
String tableName = destination.getKey().split("\\.")[1];
- LOG.info("Table Name {}", tableName);
+ LOG.debug("Table Name {}", tableName);
switch (tableName) {
case "dlp_inspection_result":
return BigQueryUtils.toTableSchema(Util.dlpInspectionSchema);
diff --git a/src/main/java/com/google/swarm/tokenization/common/DLPTransform.java b/src/main/java/com/google/swarm/tokenization/common/DLPTransform.java
index 7a23bbaf..b353cb00 100644
--- a/src/main/java/com/google/swarm/tokenization/common/DLPTransform.java
+++ b/src/main/java/com/google/swarm/tokenization/common/DLPTransform.java
@@ -169,7 +169,7 @@ public void processElement(
String deidTableName = BigQueryHelpers.parseTableSpec(element.getKey()).getTableId();
String tableName = String.format("%s_%s", deidTableName, Util.BQ_REID_TABLE_EXT);
- LOG.info("Table Ref {}", tableName);
+ LOG.debug("Table Ref {}", tableName);
Table originalData = element.getValue().getItem().getTable();
numberOfBytesReidentified.inc(originalData.toByteArray().length);
List headers =
@@ -205,7 +205,7 @@ public void processElement(
String fileName = element.getKey().split("\\~")[0];
Table tokenizedData = element.getValue().getItem().getTable();
- LOG.info("Table Tokenized {}", tokenizedData.toString());
+ LOG.debug("Table Tokenized {}", tokenizedData.toString());
numberOfRowDeidentified.inc(tokenizedData.getRowsCount());
List headers =
tokenizedData.getHeadersList().stream()
diff --git a/src/main/java/com/google/swarm/tokenization/common/Util.java b/src/main/java/com/google/swarm/tokenization/common/Util.java
index 6f61ef91..bceb3dea 100644
--- a/src/main/java/com/google/swarm/tokenization/common/Util.java
+++ b/src/main/java/com/google/swarm/tokenization/common/Util.java
@@ -90,11 +90,6 @@ public enum FileType {
new TupleTag>() {};
public static final TupleTag> customerTranscriptTuple =
new TupleTag>() {};
- public static final TupleTag> contentTag =
- new TupleTag>() {};
- public static final TupleTag> headerTag =
- new TupleTag>() {};
-
public static final TupleTag> inspectOrDeidSuccess =
new TupleTag>() {};
public static final TupleTag> inspectOrDeidFailure =
@@ -105,28 +100,15 @@ public enum FileType {
public static final TupleTag> reidFailure =
new TupleTag>() {};
- public static final TupleTag> inspectApiCallSuccess =
- new TupleTag>() {};
- public static final TupleTag> inspectApiCallError =
- new TupleTag>() {};
- public static final String BQ_DLP_INSPECT_TABLE_NAME = String.valueOf("dlp_inspection_result");
- public static final String BQ_ERROR_TABLE_NAME = String.valueOf("error_log");
- public static final String BQ_REID_TABLE_EXT = String.valueOf("re_id");
+
+ public static final String BQ_DLP_INSPECT_TABLE_NAME = "dlp_inspection_result";
+ public static final String BQ_ERROR_TABLE_NAME = "error_log";
+ public static final String BQ_REID_TABLE_EXT = "re_id";
public static final DateTimeFormatter TIMESTAMP_FORMATTER =
DateTimeFormat.forPattern("yyyy-MM-dd HH:mm:ss.SSSSSS");
- public static Table.Row convertCsvRowToTableRow(String row) {
- String[] values = row.split(",");
- Table.Row.Builder tableRowBuilder = Table.Row.newBuilder();
- for (String value : values) {
- tableRowBuilder.addValues(Value.newBuilder().setStringValue(value).build());
- }
-
- return tableRowBuilder.build();
- }
-
public static String checkHeaderName(String name) {
String checkedHeader = name.replaceAll("\\s", "_");
checkedHeader = checkedHeader.replaceAll("'", "");
@@ -181,17 +163,6 @@ public static BufferedReader getReader(ReadableFile csvFile) {
return br;
}
- public static boolean isDefaultMode(String[] args) {
-
- for (String arg : args) {
- String[] splitFromEqual = arg.split("=");
- String value = splitFromEqual[1];
- if (value.equals("s3")) {
- return false;
- }
- }
- return true;
- }
static {
DateTimeFormatter dateTimePart =
@@ -268,22 +239,8 @@ public static TableRow toTableRow(Row row) {
return output;
}
- public static List getFileHeaders(BufferedReader reader) {
- List headers = new ArrayList<>();
- try {
- CSVRecord csvHeader = CSVFormat.DEFAULT.parse(reader).getRecords().get(0);
- csvHeader.forEach(
- headerValue -> {
- headers.add(headerValue);
- });
- } catch (IOException e) {
- LOG.error("Failed to get csv header values}", e.getMessage());
- throw new RuntimeException(e);
- }
- return headers;
- }
-
@SuppressWarnings("null")
+ // custom parser for csv- will be taken out
public static List parseLine(String cvsLine, char separators, char customQuote) {
List result = new ArrayList<>();
@@ -362,7 +319,6 @@ public static List parseLine(String cvsLine, char separators, char custo
}
result.add(curVal.toString());
-
return result;
}
@@ -375,7 +331,7 @@ public static TableRow createBqRow(Table.Row tokenizedValue, String[] headers) {
.forEach(
value -> {
String checkedHeaderName =
- Util.checkHeaderName(headers[headerIndex.getAndIncrement()].toString());
+ Util.checkHeaderName(headers[headerIndex.getAndIncrement()]);
bqRow.set(checkedHeaderName, value.getStringValue());
cells.add(new TableCell().set(checkedHeaderName, value.getStringValue()));
});
@@ -423,7 +379,6 @@ public static String parseChatlog(String fileName, String chatLog) {
position = position + 1;
String finalTranscript =
Arrays.asList(tempValues).stream().collect(Collectors.joining(","));
- LOG.info(finalTranscript);
return finalTranscript;
}
@@ -442,7 +397,6 @@ public static String parseChatlog(String fileName, String chatLog) {
position = position + 1;
String finalTranscript =
Arrays.asList(tempValues).stream().collect(Collectors.joining(","));
- LOG.info(finalTranscript);
return finalTranscript;
}
@@ -461,4 +415,13 @@ public static String parseChatlog(String fileName, String chatLog) {
return StringUtils.EMPTY;
}
+ // needed to transform after parsing contextual text io read metadata for csv files.
+ public static String sanitizeCSVFileName(String fileName){
+ // fileName.csv
+ if(fileName.contains(".")){
+ return fileName.replace(".","_");
+ }
+ return fileName;
+
+ }
}