Skip to content
This repository was archived by the owner on Oct 16, 2025. It is now read-only.

Commit 4c3b61c

Browse files
committed
ready for test
1 parent 225ea3c commit 4c3b61c

7 files changed

Lines changed: 141 additions & 131 deletions

File tree

build.gradle

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ dependencies {
7979
compile group: 'org.apache.beam', name: 'beam-runners-google-cloud-dataflow-java', version: dataflowBeamVersion
8080
compile group: 'org.apache.beam', name: 'beam-runners-direct-java', version: dataflowBeamVersion
8181
compile group: 'org.slf4j', name: 'slf4j-jdk14', version: '1.7.5'
82-
compile group: 'com.google.cloud', name: 'google-cloud-dlp', version: '1.1.4'
82+
compile group: 'com.google.cloud', name: 'google-cloud-dlp', version: '1.1.4'
8383
compile group: 'com.google.apis', name: 'google-api-services-cloudkms', version: 'v1-rev108-1.25.0'
8484
compile group: 'org.apache.beam', name: 'beam-sdks-java-io-amazon-web-services', version: dataflowBeamVersion
8585
compile "com.google.auto.value:auto-value-annotations:1.6.2"

src/main/java/com/google/swarm/tokenization/DLPBqScanner.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
import java.io.IOException;
2626
import java.util.Arrays;
2727
import org.apache.avro.Schema;
28+
import org.apache.beam.runners.dataflow.options.DataflowPipelineDebugOptions;
2829
import org.apache.beam.sdk.Pipeline;
2930
import org.apache.beam.sdk.PipelineResult;
3031
import org.apache.beam.sdk.coders.AvroCoder;
@@ -46,7 +47,8 @@ public class DLPBqScanner {
4647
public static void main(String[] args) {
4748
DLPBqScannerOptions options =
4849
PipelineOptionsFactory.fromArgs(args).withValidation().as(DLPBqScannerOptions.class);
49-
50+
options.setDumpHeapOnOOM(true);
51+
options.setSaveHeapDumpsToGcsPath("gs://storage-api-hd");
5052
run(options);
5153
}
5254

@@ -62,7 +64,7 @@ static BigQueryStorageClient create() {
6264
}
6365
}
6466

65-
public interface DLPBqScannerOptions extends PipelineOptions {
67+
public interface DLPBqScannerOptions extends PipelineOptions, DataflowPipelineDebugOptions {
6668
@Description("BigQuery table to export from in the form <project>:<dataset>.<table>")
6769
@Required
6870
String getTableRef();

src/main/java/com/google/swarm/tokenization/DLPS3ScannerPipeline.java

Lines changed: 46 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -15,29 +15,15 @@
1515
*/
1616
package com.google.swarm.tokenization;
1717

18-
import com.google.swarm.tokenization.common.AuditInspectDataTransform;
1918
import com.google.swarm.tokenization.common.DLPTransform;
2019
import com.google.swarm.tokenization.common.FileReaderTransform;
21-
import com.google.swarm.tokenization.common.RowToJson;
2220
import com.google.swarm.tokenization.common.S3ReaderOptions;
23-
import com.google.swarm.tokenization.common.Util;
2421
import org.apache.beam.sdk.Pipeline;
2522
import org.apache.beam.sdk.PipelineResult;
26-
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO;
27-
import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO;
2823
import org.apache.beam.sdk.options.PipelineOptionsFactory;
29-
import org.apache.beam.sdk.transforms.DoFn;
30-
import org.apache.beam.sdk.transforms.DoFn.ProcessContext;
31-
import org.apache.beam.sdk.transforms.DoFn.ProcessElement;
32-
import org.apache.beam.sdk.transforms.ParDo;
33-
import org.apache.beam.sdk.transforms.windowing.AfterProcessingTime;
34-
import org.apache.beam.sdk.transforms.windowing.FixedWindows;
35-
import org.apache.beam.sdk.transforms.windowing.Window;
3624
import org.apache.beam.sdk.values.KV;
3725
import org.apache.beam.sdk.values.PCollection;
3826
import org.apache.beam.sdk.values.PCollectionTuple;
39-
import org.apache.beam.sdk.values.Row;
40-
import org.joda.time.Duration;
4127
import org.slf4j.Logger;
4228
import org.slf4j.LoggerFactory;
4329

@@ -56,28 +42,29 @@ public static PipelineResult run(S3ReaderOptions options) {
5642

5743
PCollection<KV<String, String>> nonInspectedContents =
5844
p.apply(
59-
"File Read Transform",
60-
FileReaderTransform.newBuilder()
61-
.setSubscriber(options.getSubscriber())
62-
.setDelimeter(options.getDelimeter())
63-
.setKeyRange(options.getKeyRange())
64-
.build());
65-
// .apply(
66-
// "Fixed Window",
67-
// Window.<KV<String, String>>into(FixedWindows.of(Duration.standardSeconds(10)))
68-
// .triggering(
69-
// AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.ZERO))
70-
// .discardingFiredPanes()
71-
// .withAllowedLateness(Duration.ZERO));
72-
73-
// nonInspectedContents.apply("Print", ParDo.of(new DoFn<KV<String,String>, String>(){
74-
//
75-
// @ProcessElement
76-
// public void processElement(ProcessContext c) {
77-
// c.output(c.element().getValue());
78-
// }
79-
// }));
45+
"File Read Transform",
46+
FileReaderTransform.newBuilder()
47+
.setSubscriber(options.getSubscriber())
48+
.setDelimeter(options.getDelimeter())
49+
.setKeyRange(options.getKeyRange())
50+
.build());
51+
// .apply(
52+
// "Fixed Window",
53+
// Window.<KV<String, String>>into(FixedWindows.of(Duration.standardSeconds(10)))
54+
// .triggering(
55+
//
56+
// AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.ZERO))
57+
// .discardingFiredPanes()
58+
// .withAllowedLateness(Duration.ZERO));
8059

60+
// nonInspectedContents.apply("Print", ParDo.of(new DoFn<KV<String,String>, String>(){
61+
//
62+
// @ProcessElement
63+
// public void processElement(ProcessContext c) {
64+
// c.output(c.element().getValue());
65+
// }
66+
// }));
67+
//
8168
PCollectionTuple inspectedData =
8269
nonInspectedContents.apply(
8370
"DLPScanner",
@@ -87,30 +74,30 @@ public static PipelineResult run(S3ReaderOptions options) {
8774
.setBatchSize(options.getBatchSize())
8875
.build());
8976

90-
// PCollection<Row> inspectedContents =
91-
// inspectedData.get(Util.inspectData).setRowSchema(Util.bqDataSchema);
92-
//
93-
// PCollection<Row> inspectedStats =
94-
// inspectedData.get(Util.auditData).setRowSchema(Util.bqAuditSchema);
95-
//
96-
// PCollection<Row> auditData =
97-
// inspectedStats
98-
// .apply("FileTrackerTransform", new AuditInspectDataTransform())
99-
// .setRowSchema(Util.bqAuditSchema);
100-
//
101-
// auditData.apply(
102-
// "WriteAuditData",
103-
// BigQueryIO.<Row>write()
104-
// .to(options.getAuditTableSpec())
105-
// .withMethod(BigQueryIO.Write.Method.STREAMING_INSERTS)
106-
// .useBeamSchema()
107-
// .withoutValidation()
108-
// .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
109-
// .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER));
110-
//
111-
// auditData
112-
// .apply("RowToJson", new RowToJson())
113-
// .apply("WriteToTopic", PubsubIO.writeStrings().to(options.getTopic()));
77+
// PCollection<Row> inspectedContents =
78+
// inspectedData.get(Util.inspectData).setRowSchema(Util.bqDataSchema);
79+
//
80+
// PCollection<Row> inspectedStats =
81+
// inspectedData.get(Util.auditData).setRowSchema(Util.bqAuditSchema);
82+
//
83+
// PCollection<Row> auditData =
84+
// inspectedStats
85+
// .apply("FileTrackerTransform", new AuditInspectDataTransform())
86+
// .setRowSchema(Util.bqAuditSchema);
87+
//
88+
// auditData.apply(
89+
// "WriteAuditData",
90+
// BigQueryIO.<Row>write()
91+
// .to(options.getAuditTableSpec())
92+
// .withMethod(BigQueryIO.Write.Method.STREAMING_INSERTS)
93+
// .useBeamSchema()
94+
// .withoutValidation()
95+
// .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
96+
// .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER));
97+
//
98+
// auditData
99+
// .apply("RowToJson", new RowToJson())
100+
// .apply("WriteToTopic", PubsubIO.writeStrings().to(options.getTopic()));
114101

115102
// inspectedContents.apply(
116103
// "WriteInspectData",

src/main/java/com/google/swarm/tokenization/common/DLPTransform.java

Lines changed: 84 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,10 @@
1515
*/
1616
package com.google.swarm.tokenization.common;
1717

18+
import com.google.api.gax.retrying.RetrySettings;
1819
import com.google.auto.value.AutoValue;
1920
import com.google.cloud.dlp.v2.DlpServiceClient;
21+
import com.google.cloud.dlp.v2.DlpServiceSettings;
2022
import com.google.privacy.dlp.v2.ContentItem;
2123
import com.google.privacy.dlp.v2.InspectContentRequest;
2224
import com.google.privacy.dlp.v2.InspectContentResponse;
@@ -43,9 +45,9 @@
4345
import org.apache.beam.sdk.values.PCollectionTuple;
4446
import org.apache.beam.sdk.values.Row;
4547
import org.apache.beam.sdk.values.TupleTagList;
46-
import org.joda.time.Duration;
4748
import org.slf4j.Logger;
4849
import org.slf4j.LoggerFactory;
50+
import org.threeten.bp.Duration;
4951

5052
@AutoValue
5153
public abstract class DLPTransform
@@ -88,8 +90,8 @@ public PCollectionTuple expand(PCollection<KV<String, String>> input) {
8890
public static class BatchRequest extends DoFn<KV<String, String>, KV<String, Iterable<String>>> {
8991

9092
private static final long serialVersionUID = 1L;
91-
// private final Counter numberOfRowsBagged =
92-
// Metrics.counter(BatchRequest.class, "numberOfRowsBagged");
93+
// private final Counter numberOfRowsBagged =
94+
// Metrics.counter(BatchRequest.class, "numberOfRowsBagged");
9395
private Integer batchSize;
9496

9597
public BatchRequest(Integer batchSize) {
@@ -110,7 +112,7 @@ public void process(
110112
@StateId("elementsBag") BagState<KV<String, String>> elementsBag) {
111113
elementsBag.add(element);
112114
// eventTimer.set(w.maxTimestamp());
113-
eventTimer.offset(Duration.standardSeconds(30)).setRelative();
115+
eventTimer.offset(org.joda.time.Duration.standardSeconds(10)).setRelative();
114116
}
115117

116118
@OnTimer("eventTimer")
@@ -127,7 +129,7 @@ public void onTimer(
127129
Integer elementSize = element.getValue().getBytes().length;
128130
boolean clearBuffer = bufferSize.intValue() + elementSize.intValue() > batchSize;
129131
if (clearBuffer) {
130-
//numberOfRowsBagged.inc(rows.size());
132+
// numberOfRowsBagged.inc(rows.size());
131133
LOG.debug("Clear Buffer {} , Key {}", bufferSize.intValue(), key);
132134
output.output(KV.of(key, rows));
133135
// clean up in a method
@@ -145,7 +147,7 @@ public void onTimer(
145147
// must be a better way
146148
if (!rows.isEmpty()) {
147149
LOG.debug("Remaining buffer {}, key{}", rows.size(), key);
148-
//numberOfRowsBagged.inc(rows.size());
150+
// numberOfRowsBagged.inc(rows.size());
149151
output.output(KV.of(key, rows));
150152
}
151153
}
@@ -157,6 +159,7 @@ public static class InspectData extends DoFn<KV<String, Iterable<String>>, Row>
157159
private InspectContentRequest.Builder requestBuilder;
158160
private final Counter numberOfBytesInspected =
159161
Metrics.counter(InspectData.class, "NumberOfBytesInspected");
162+
private DlpServiceClient dlpServiceClient;
160163

161164
public InspectData(String projectId, String inspectTemplateName) {
162165
this.projectId = projectId;
@@ -171,63 +174,83 @@ public void setup() {
171174
.setInspectTemplateName(this.inspectTemplateName);
172175
}
173176

177+
@StartBundle
178+
public void startBundle() throws IOException {
179+
180+
// DlpServiceSettings.Builder settingsBuilder = DlpServiceSettings.newBuilder();
181+
// settingsBuilder
182+
// .inspectContentSettings()
183+
// .setRetrySettings(
184+
// RetrySettings.newBuilder()
185+
// .setInitialRpcTimeout(Duration.ofSeconds(60))
186+
// .setMaxRpcTimeout(Duration.ofSeconds(60))
187+
// .build());
188+
dlpServiceClient = DlpServiceClient.create();
189+
}
190+
191+
@FinishBundle
192+
public void finishBundle() {
193+
if (dlpServiceClient != null) {
194+
dlpServiceClient.close();
195+
}
196+
}
197+
174198
@ProcessElement
175199
public void processElement(ProcessContext c) throws IOException {
176-
try (DlpServiceClient dlpServiceClient = DlpServiceClient.create()) {
177-
String fileName = c.element().getKey();
178-
ContentItem contentItem =
179-
ContentItem.newBuilder().setValue(emitResult(c.element().getValue())).build();
180-
this.requestBuilder.setItem(contentItem);
181-
InspectContentResponse response =
182-
dlpServiceClient.inspectContent(this.requestBuilder.build());
183-
String timeStamp = Util.getTimeStamp();
184-
185-
boolean hasErrors = response.findInitializationErrors().stream().count() > 0;
186-
if (response.hasResult() && !hasErrors) {
187-
long bytesInspected = contentItem.getValue().getBytes().length;
188-
response
189-
.getResult()
190-
.getFindingsList()
191-
.forEach(
192-
finding -> {
193-
Row row =
194-
Row.withSchema(Util.bqDataSchema)
195-
.addValues(
196-
fileName,
197-
timeStamp,
198-
finding.getInfoType().getName(),
199-
finding.getLikelihood().name(),
200-
finding.getQuote(),
201-
finding.getLocation().getCodepointRange().getStart(),
202-
finding.getLocation().getCodepointRange().getEnd())
203-
.build();
204-
LOG.debug("Row {}", row);
205-
206-
c.output(Util.inspectData, row);
207-
});
208-
numberOfBytesInspected.inc(bytesInspected);
209-
c.output(
210-
Util.auditData,
211-
Row.withSchema(Util.bqAuditSchema)
212-
.addValues(fileName, timeStamp, bytesInspected, Util.INSPECTED)
213-
.build());
214-
} else {
215-
response
216-
.findInitializationErrors()
217-
.forEach(
218-
error -> {
219-
c.output(
220-
Util.errorData,
221-
Row.withSchema(Util.errorSchema)
222-
.addValues(fileName, timeStamp, error.toString())
223-
.build());
224-
});
225-
c.output(
226-
Util.auditData,
227-
Row.withSchema(Util.bqAuditSchema)
228-
.addValues(fileName, timeStamp, 0, Util.FAILED)
229-
.build());
230-
}
200+
201+
String fileName = c.element().getKey();
202+
String contents = emitResult(c.element().getValue());
203+
ContentItem contentItem = ContentItem.newBuilder().setValue(contents).build();
204+
this.requestBuilder.setItem(contentItem);
205+
InspectContentResponse response =
206+
dlpServiceClient.inspectContent(this.requestBuilder.build());
207+
String timeStamp = Util.getTimeStamp();
208+
209+
boolean hasErrors = response.findInitializationErrors().stream().count() > 0;
210+
if (response.hasResult() && !hasErrors) {
211+
long bytesInspected = contents.getBytes().length;
212+
response
213+
.getResult()
214+
.getFindingsList()
215+
.forEach(
216+
finding -> {
217+
Row row =
218+
Row.withSchema(Util.bqDataSchema)
219+
.addValues(
220+
fileName,
221+
timeStamp,
222+
finding.getInfoType().getName(),
223+
finding.getLikelihood().name(),
224+
finding.getQuote(),
225+
finding.getLocation().getCodepointRange().getStart(),
226+
finding.getLocation().getCodepointRange().getEnd())
227+
.build();
228+
LOG.debug("Row {}", row);
229+
230+
c.output(Util.inspectData, row);
231+
});
232+
numberOfBytesInspected.inc(bytesInspected);
233+
c.output(
234+
Util.auditData,
235+
Row.withSchema(Util.bqAuditSchema)
236+
.addValues(fileName, timeStamp, bytesInspected, Util.INSPECTED)
237+
.build());
238+
} else {
239+
response
240+
.findInitializationErrors()
241+
.forEach(
242+
error -> {
243+
c.output(
244+
Util.errorData,
245+
Row.withSchema(Util.errorSchema)
246+
.addValues(fileName, timeStamp, error.toString())
247+
.build());
248+
});
249+
c.output(
250+
Util.auditData,
251+
Row.withSchema(Util.bqAuditSchema)
252+
.addValues(fileName, timeStamp, 0, Util.FAILED)
253+
.build());
231254
}
232255
}
233256
}

0 commit comments

Comments
 (0)