1515 */
1616package com .google .swarm .tokenization ;
1717
18+ import com .google .swarm .tokenization .common .AuditInspectDataTransform ;
19+ import com .google .swarm .tokenization .common .BQWriteTransform ;
1820import com .google .swarm .tokenization .common .DLPTransform ;
1921import com .google .swarm .tokenization .common .FileReaderTransform ;
22+ import com .google .swarm .tokenization .common .RowToJson ;
2023import com .google .swarm .tokenization .common .S3ReaderOptions ;
24+ import com .google .swarm .tokenization .common .Util ;
2125import org .apache .beam .sdk .Pipeline ;
2226import org .apache .beam .sdk .PipelineResult ;
27+ import org .apache .beam .sdk .io .gcp .bigquery .BigQueryIO ;
28+ import org .apache .beam .sdk .io .gcp .pubsub .PubsubIO ;
2329import org .apache .beam .sdk .options .PipelineOptionsFactory ;
2430import org .apache .beam .sdk .values .KV ;
2531import org .apache .beam .sdk .values .PCollection ;
2632import org .apache .beam .sdk .values .PCollectionTuple ;
33+ import org .apache .beam .sdk .values .Row ;
2734import org .slf4j .Logger ;
2835import org .slf4j .LoggerFactory ;
2936
@@ -33,7 +40,6 @@ public class DLPS3ScannerPipeline {
3340 public static void main (String [] args ) {
3441 S3ReaderOptions options =
3542 PipelineOptionsFactory .fromArgs (args ).withValidation ().as (S3ReaderOptions .class );
36- // options.setEnableStreamingEngine(true);
3743 run (options );
3844 }
3945
@@ -48,23 +54,7 @@ public static PipelineResult run(S3ReaderOptions options) {
4854 .setDelimeter (options .getDelimeter ())
4955 .setKeyRange (options .getKeyRange ())
5056 .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));
5957
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- //
6858 PCollectionTuple inspectedData =
6959 nonInspectedContents .apply (
7060 "DLPScanner" ,
@@ -74,37 +64,37 @@ public static PipelineResult run(S3ReaderOptions options) {
7464 .setBatchSize (options .getBatchSize ())
7565 .build ());
7666
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()));
67+ PCollection <Row > inspectedContents =
68+ inspectedData .get (Util .inspectData ).setRowSchema (Util .bqDataSchema );
10169
102- // inspectedContents.apply(
103- // "WriteInspectData",
104- // BQWriteTransform.newBuilder()
105- // .setTableSpec(options.getTableSpec())
106- // .setMethod(options.getWriteMethod())
107- // .build());
70+ PCollection <Row > inspectedStats =
71+ inspectedData .get (Util .auditData ).setRowSchema (Util .bqAuditSchema );
72+
73+ PCollection <Row > auditData =
74+ inspectedStats
75+ .apply ("FileTrackerTransform" , new AuditInspectDataTransform ())
76+ .setRowSchema (Util .bqAuditSchema );
77+
78+ auditData .apply (
79+ "WriteAuditData" ,
80+ BigQueryIO .<Row >write ()
81+ .to (options .getAuditTableSpec ())
82+ .withMethod (BigQueryIO .Write .Method .STREAMING_INSERTS )
83+ .useBeamSchema ()
84+ .withoutValidation ()
85+ .withWriteDisposition (BigQueryIO .Write .WriteDisposition .WRITE_APPEND )
86+ .withCreateDisposition (BigQueryIO .Write .CreateDisposition .CREATE_NEVER ));
87+
88+ auditData
89+ .apply ("RowToJson" , new RowToJson ())
90+ .apply ("WriteToTopic" , PubsubIO .writeStrings ().to (options .getTopic ()));
91+
92+ inspectedContents .apply (
93+ "WriteInspectData" ,
94+ BQWriteTransform .newBuilder ()
95+ .setTableSpec (options .getTableSpec ())
96+ .setMethod (options .getWriteMethod ())
97+ .build ());
10898 return p .run ();
10999 }
110100}
0 commit comments