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; + + } }