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

Commit 225ea3c

Browse files
committed
byte change mismatch but completed in 30 mins
1 parent 01a1b03 commit 225ea3c

3 files changed

Lines changed: 64 additions & 54 deletions

File tree

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

Lines changed: 43 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,10 @@
2626
import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO;
2727
import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO;
2828
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;
2933
import org.apache.beam.sdk.transforms.windowing.AfterProcessingTime;
3034
import org.apache.beam.sdk.transforms.windowing.FixedWindows;
3135
import org.apache.beam.sdk.transforms.windowing.Window;
@@ -57,22 +61,22 @@ public static PipelineResult run(S3ReaderOptions options) {
5761
.setSubscriber(options.getSubscriber())
5862
.setDelimeter(options.getDelimeter())
5963
.setKeyRange(options.getKeyRange())
60-
.build())
61-
.apply(
62-
"Fixed Window",
63-
Window.<KV<String, String>>into(FixedWindows.of(Duration.standardSeconds(1)))
64-
.triggering(
65-
AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.ZERO))
66-
.discardingFiredPanes()
67-
.withAllowedLateness(Duration.ZERO));
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));
6872

69-
// nonInspectedContents.apply("Print", ParDo.of(new DoFn<KV<String,String>, String>(){
70-
//
71-
// @ProcessElement
72-
// public void processElement(ProcessContext c) {
73-
// c.output(c.element().getValue());
74-
// }
75-
// }));
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+
// }));
7680

7781
PCollectionTuple inspectedData =
7882
nonInspectedContents.apply(
@@ -83,30 +87,30 @@ public static PipelineResult run(S3ReaderOptions options) {
8387
.setBatchSize(options.getBatchSize())
8488
.build());
8589

86-
PCollection<Row> inspectedContents =
87-
inspectedData.get(Util.inspectData).setRowSchema(Util.bqDataSchema);
88-
89-
PCollection<Row> inspectedStats =
90-
inspectedData.get(Util.auditData).setRowSchema(Util.bqAuditSchema);
91-
92-
PCollection<Row> auditData =
93-
inspectedStats
94-
.apply("FileTrackerTransform", new AuditInspectDataTransform())
95-
.setRowSchema(Util.bqAuditSchema);
96-
97-
auditData.apply(
98-
"WriteAuditData",
99-
BigQueryIO.<Row>write()
100-
.to(options.getAuditTableSpec())
101-
.withMethod(BigQueryIO.Write.Method.STREAMING_INSERTS)
102-
.useBeamSchema()
103-
.withoutValidation()
104-
.withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
105-
.withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER));
106-
107-
auditData
108-
.apply("RowToJson", new RowToJson())
109-
.apply("WriteToTopic", PubsubIO.writeStrings().to(options.getTopic()));
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()));
110114

111115
// inspectedContents.apply(
112116
// "WriteInspectData",

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

Lines changed: 10 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -88,8 +88,8 @@ public PCollectionTuple expand(PCollection<KV<String, String>> input) {
8888
public static class BatchRequest extends DoFn<KV<String, String>, KV<String, Iterable<String>>> {
8989

9090
private static final long serialVersionUID = 1L;
91-
private final Counter numberOfRowsBagged =
92-
Metrics.counter(BatchRequest.class, "numberOfRowsBagged");
91+
// private final Counter numberOfRowsBagged =
92+
// Metrics.counter(BatchRequest.class, "numberOfRowsBagged");
9393
private Integer batchSize;
9494

9595
public BatchRequest(Integer batchSize) {
@@ -110,7 +110,7 @@ public void process(
110110
@StateId("elementsBag") BagState<KV<String, String>> elementsBag) {
111111
elementsBag.add(element);
112112
// eventTimer.set(w.maxTimestamp());
113-
eventTimer.offset(Duration.standardSeconds(5)).setRelative();
113+
eventTimer.offset(Duration.standardSeconds(30)).setRelative();
114114
}
115115

116116
@OnTimer("eventTimer")
@@ -127,8 +127,8 @@ public void onTimer(
127127
Integer elementSize = element.getValue().getBytes().length;
128128
boolean clearBuffer = bufferSize.intValue() + elementSize.intValue() > batchSize;
129129
if (clearBuffer) {
130-
numberOfRowsBagged.inc(rows.size());
131-
LOG.info("Clear Buffer {} , Key {}", bufferSize.intValue(), key);
130+
//numberOfRowsBagged.inc(rows.size());
131+
LOG.debug("Clear Buffer {} , Key {}", bufferSize.intValue(), key);
132132
output.output(KV.of(key, rows));
133133
// clean up in a method
134134
rows.clear();
@@ -144,8 +144,8 @@ public void onTimer(
144144
});
145145
// must be a better way
146146
if (!rows.isEmpty()) {
147-
LOG.info("Remaining buffer {}, key{}", rows.size(), key);
148-
numberOfRowsBagged.inc(rows.size());
147+
LOG.debug("Remaining buffer {}, key{}", rows.size(), key);
148+
//numberOfRowsBagged.inc(rows.size());
149149
output.output(KV.of(key, rows));
150150
}
151151
}
@@ -181,14 +181,11 @@ public void processElement(ProcessContext c) throws IOException {
181181
InspectContentResponse response =
182182
dlpServiceClient.inspectContent(this.requestBuilder.build());
183183
String timeStamp = Util.getTimeStamp();
184-
185-
long bytesInspected = contentItem.getSerializedSize();
186-
int totalFinding =
187-
Long.valueOf(response.getResult().getFindingsList().stream().count()).intValue();
188-
LOG.debug("bytes inspected {}", bytesInspected);
184+
189185
boolean hasErrors = response.findInitializationErrors().stream().count() > 0;
190186
if (response.hasResult() && !hasErrors) {
191-
response
187+
long bytesInspected = contentItem.getValue().getBytes().length;
188+
response
192189
.getResult()
193190
.getFindingsList()
194191
.forEach(

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

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
import java.util.Random;
2222
import org.apache.beam.sdk.io.FileIO.ReadableFile;
2323
import org.apache.beam.sdk.io.range.OffsetRange;
24+
import org.apache.beam.sdk.metrics.Counter;
25+
import org.apache.beam.sdk.metrics.Metrics;
2426
import org.apache.beam.sdk.transforms.DoFn;
2527
import org.apache.beam.sdk.transforms.splittabledofn.OffsetRangeTracker;
2628
import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker;
@@ -29,12 +31,17 @@
2931
import org.slf4j.Logger;
3032
import org.slf4j.LoggerFactory;
3133

34+
import com.google.swarm.tokenization.common.DLPTransform.BatchRequest;
35+
3236
public class FileReaderSplitDoFn extends DoFn<KV<String, ReadableFile>, KV<String, String>> {
3337
public static final Logger LOG = LoggerFactory.getLogger(FileReaderSplitDoFn.class);
3438
public static Integer SPLIT_SIZE = 1000000;
3539
private String delimeter;
3640
private Integer keyRange;
37-
41+
private final Counter numberOfRows =
42+
Metrics.counter(FileReaderSplitDoFn.class, "numberOfRows");
43+
private final Counter numberOfBytesRead =
44+
Metrics.counter(FileReaderSplitDoFn.class, "numberOfBytesRead");
3845
public FileReaderSplitDoFn(String delimeter, Integer keyRange) {
3946
this.delimeter = delimeter;
4047
this.keyRange = keyRange;
@@ -51,7 +58,9 @@ public void processElement(ProcessContext c, RestrictionTracker<OffsetRange, Lon
5158
reader.readNextRecord();
5259
String contents = reader.getCurrent();
5360
String key = String.format("%s~%d", fileName, new Random().nextInt(keyRange));
54-
c.outputWithTimestamp(KV.of(key, contents), Instant.now());
61+
numberOfRows.inc();
62+
numberOfBytesRead.inc(contents.length());
63+
c.output(KV.of(key, contents));
5564
}
5665
}
5766
}

0 commit comments

Comments
 (0)