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

Commit f93cc83

Browse files
authored
Merge pull request #157 from chitara-01/parquet-support
Support to INSPECT & DEID parquet files
2 parents aa5c2f6 + 8829537 commit f93cc83

20 files changed

Lines changed: 717 additions & 22 deletions
13.8 MB
Binary file not shown.
13.8 MB
Binary file not shown.
13.8 MB
Binary file not shown.
13.8 MB
Binary file not shown.
13.8 MB
Binary file not shown.

.github/workflows/dlp-pipelines.yml

Lines changed: 196 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -9,17 +9,23 @@ on:
99
env:
1010
PROJECT_ID: "dlp-dataflow-deid-ci-392604"
1111
DATASET_ID: "demo_dataset"
12+
PARQUET_DATASET_ID: "parquet_results"
13+
REGION: "us-central1"
1214
GCS_BUCKET: "dlp-dataflow-deid-ci-392604-demo-data"
1315
GCS_NOTIFICATION_TOPIC: "projects/dlp-dataflow-deid-ci-392604/topics/dlp-dataflow-deid-ci-gcs-notification-topic"
16+
SAMPLE_DATA_DIR: "sample_data_for_ci_workflow"
1417
INPUT_FILE_NAME: "tiny_csv"
1518
INPUT_STREAMING_WRITE_FILE_NAME: "streaming_write"
19+
INPUT_PARQUET_FILE_NAME: "tiny_parquet"
1620
INSPECT_TEMPLATE_PATH: "projects/dlp-dataflow-deid-ci-392604/locations/global/inspectTemplates/dlp-demo-inspect-latest-1689137435622"
1721
DEID_TEMPLATE_PATH: "projects/dlp-dataflow-deid-ci-392604/locations/global/deidentifyTemplates/dlp-demo-deid-latest-1689137435622"
22+
PARQUET_DEID_TEMPLATE_PATH: "projects/dlp-dataflow-deid-ci-392604/locations/global/deidentifyTemplates/parquet-dlp-demo-deid-latest-1689137435622"
1823
REID_TEMPLATE_PATH: "projects/dlp-dataflow-deid-ci-392604/locations/global/deidentifyTemplates/dlp-demo-reid-latest-1689137435622"
1924
SERVICE_ACCOUNT_EMAIL: "demo-service-account@dlp-dataflow-deid-ci-392604.iam.gserviceaccount.com"
2025
PUBSUB_TOPIC_NAME: "demo-topic"
2126
INSPECTION_TABLE_ID: "dlp_inspection_result"
2227
NUM_INSPECTION_RECORDS_THRESHOLD: "50"
28+
PARQUET_INSPECTION_RECORDS_THRESHOLD: "30"
2329
REIDENTIFICATION_QUERY_FILE: "reid_query.sql"
2430

2531
jobs:
@@ -30,10 +36,13 @@ jobs:
3036

3137
outputs:
3238
output1: ${{ steps.gen-uuid.outputs.uuid }}
33-
output2: ${{ steps.gen-uuid.outputs.inpect_job_name }}
39+
output2: ${{ steps.gen-uuid.outputs.inspect_job_name }}
3440
output3: ${{ steps.gen-uuid.outputs.deid_job_name }}
3541
output4: ${{ steps.gen-uuid.outputs.reid_job_name }}
3642
output5: ${{ steps.gen-uuid.outputs.dataset_id }}
43+
output6: ${{ steps.gen-uuid.outputs.parquet_inspect_job_name }}
44+
output7: ${{ steps.gen-uuid.outputs.parquet_deid_job_name }}
45+
output8: ${{ steps.gen-uuid.outputs.parquet_dataset_id }}
3746

3847
steps:
3948
- name: Generate UUID for workflow
@@ -42,10 +51,13 @@ jobs:
4251
new_uuid=$(uuidgen)
4352
modified_uuid=$(echo "$new_uuid" | tr '-' '_')
4453
echo "uuid=$new_uuid" >> "$GITHUB_OUTPUT"
45-
echo "inpect_job_name=inspect-$new_uuid" >> "$GITHUB_OUTPUT"
54+
echo "inspect_job_name=inspect-$new_uuid" >> "$GITHUB_OUTPUT"
4655
echo "deid_job_name=deid-$new_uuid" >> "$GITHUB_OUTPUT"
4756
echo "reid_job_name=reid-$new_uuid" >> "$GITHUB_OUTPUT"
4857
echo "dataset_id=${{ env.DATASET_ID }}_$modified_uuid" >> "$GITHUB_OUTPUT"
58+
echo "parquet_inspect_job_name=parquet-inspect-$new_uuid" >> "$GITHUB_OUTPUT"
59+
echo "parquet_deid_job_name=parquet-deid-$new_uuid" >> "$GITHUB_OUTPUT"
60+
echo "parquet_dataset_id=${{ env.PARQUET_DATASET_ID }}_$modified_uuid" >> "$GITHUB_OUTPUT"
4961
5062
create-dataset:
5163
needs:
@@ -60,8 +72,10 @@ jobs:
6072
- name: Create BQ dataset
6173
env:
6274
DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }}
75+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
6376
run: |
6477
bq --location=US mk -d --description "GitHub CI workflow dataset" ${{ env.DATASET_ID }}
78+
bq --location=US mk -d --description "GitHub CI workflow dataset to store parquet results" ${{ env.PARQUET_DATASET_ID }}
6579
6680
inspection:
6781
needs:
@@ -91,7 +105,7 @@ jobs:
91105
run: |
92106
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
93107
--streaming --enableStreamingEngine \
94-
--region=us-central1 \
108+
--region=${{env.REGION}} \
95109
--project=${{env.PROJECT_ID}} \
96110
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
97111
--numWorkers=1 \
@@ -103,7 +117,7 @@ jobs:
103117
--inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \
104118
--batchSize=200000 \
105119
--DLPMethod=INSPECT \
106-
--serviceAccount=$SERVICE_ACCOUNT_EMAIL \
120+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
107121
--jobName=${{env.INSPECT_JOB_NAME}} \
108122
--gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}"
109123
sleep 30s
@@ -146,6 +160,84 @@ jobs:
146160
done
147161
echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count."
148162
163+
# Inspect only existing parquet files
164+
inspect-parquet-data:
165+
needs:
166+
- generate-uuid
167+
- create-dataset
168+
169+
runs-on:
170+
- self-hosted
171+
172+
timeout-minutes: 30
173+
174+
steps:
175+
- uses: actions/checkout@v2
176+
177+
- name: Setup Java
178+
uses: actions/setup-java@v1
179+
with:
180+
java-version: 11
181+
182+
- name: Setup Gradle
183+
uses: gradle/gradle-build-action@v2
184+
185+
- name: Run DLP Pipeline
186+
env:
187+
PARQUET_INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output6 }}
188+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
189+
run: |
190+
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
191+
--region=${{env.REGION}} \
192+
--project=${{env.PROJECT_ID}} \
193+
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
194+
--numWorkers=1 --maxNumWorkers=2 \
195+
--runner=DataflowRunner \
196+
--filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}*.parquet \
197+
--dataset=${{env.PARQUET_DATASET_ID}} \
198+
--inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \
199+
--batchSize=200000 \
200+
--DLPMethod=INSPECT \
201+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
202+
--jobName=${{env.PARQUET_INSPECT_JOB_NAME}}"
203+
204+
- name: Verify BQ table
205+
env:
206+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
207+
run: |
208+
not_verified=true
209+
table_count=0
210+
while $not_verified; do
211+
table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.PARQUET_DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INSPECTION_TABLE_ID}}"' | wc -l ) -1))
212+
if [[ "$table_count" == "1" ]]; then
213+
echo "PASSED";
214+
not_verified=false;
215+
else
216+
sleep 30s
217+
fi
218+
echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ."
219+
done
220+
echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ."
221+
222+
- name: Verify distinct rows
223+
env:
224+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
225+
run: |
226+
not_verified=true
227+
row_count=0
228+
while $not_verified; do
229+
row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(*) FROM `${{env.PROJECT_ID}}.${{env.PARQUET_DATASET_ID}}.${{env.INSPECTION_TABLE_ID}}`')
230+
row_count=$(echo "$row_count_json" | jq -r '.[].f0_')
231+
if [[ "$row_count" == ${{env.PARQUET_INSPECTION_RECORDS_THRESHOLD}} ]]; then
232+
echo "PASSED";
233+
not_verified=false;
234+
else
235+
sleep 30s
236+
fi
237+
echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count."
238+
done
239+
echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count."
240+
149241
de-identification:
150242
needs:
151243
- generate-uuid
@@ -174,7 +266,7 @@ jobs:
174266
run: |
175267
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
176268
--streaming --enableStreamingEngine \
177-
--region=us-central1 \
269+
--region=${{env.REGION}} \
178270
--project=${{env.PROJECT_ID}} \
179271
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
180272
--numWorkers=2 \
@@ -187,7 +279,7 @@ jobs:
187279
--deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \
188280
--batchSize=200000 \
189281
--DLPMethod=DEID \
190-
--serviceAccount=${SERVICE_ACCOUNT_EMAIL} \
282+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
191283
--jobName=${{env.DEID_JOB_NAME}} \
192284
--gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}"
193285
sleep 30s
@@ -255,7 +347,88 @@ jobs:
255347
done
256348
echo "# records in input CSV file are: $rc_orig."
257349
echo "Verified number of rows in ${{env.INPUT_FILE_NAME}}_pub_sub: $row_count."
258-
350+
351+
# Deidentify only existing parquet files
352+
deidentify-parquet-data:
353+
needs:
354+
- generate-uuid
355+
- create-dataset
356+
357+
runs-on:
358+
- self-hosted
359+
360+
timeout-minutes: 30
361+
362+
steps:
363+
- uses: actions/checkout@v2
364+
365+
- name: Setup Java
366+
uses: actions/setup-java@v1
367+
with:
368+
java-version: 11
369+
370+
- name: Setup Gradle
371+
uses: gradle/gradle-build-action@v2
372+
373+
- name: Run DLP Pipeline
374+
env:
375+
PARQUET_DEID_JOB_NAME: ${{ needs.generate-uuid.outputs.output7 }}
376+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
377+
run: |
378+
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
379+
--region=${{env.REGION}} \
380+
--project=${{env.PROJECT_ID}} \
381+
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
382+
--numWorkers=1 --maxNumWorkers=2 \
383+
--runner=DataflowRunner \
384+
--filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}.parquet \
385+
--dataset=${{env.PARQUET_DATASET_ID}} \
386+
--inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \
387+
--deidentifyTemplateName=${{env.PARQUET_DEID_TEMPLATE_PATH}} \
388+
--batchSize=200000 \
389+
--DLPMethod=DEID \
390+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
391+
--jobName=${{env.PARQUET_DEID_JOB_NAME}}"
392+
393+
- name: Verify BQ tables
394+
env:
395+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
396+
run: |
397+
not_verified=true
398+
table_count=0
399+
while $not_verified; do
400+
table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.PARQUET_DATASET_ID}}`.__TABLES__ WHERE table_id LIKE "${{env.INPUT_PARQUET_FILE_NAME}}%"' | wc -l ) -1))
401+
if [[ "$table_count" == "1" ]]; then
402+
echo "PASSED";
403+
not_verified=false;
404+
else
405+
sleep 30s
406+
fi
407+
echo "Got number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ."
408+
done
409+
echo "Verified number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ."
410+
411+
- name: Verify distinct rows of existing file
412+
env:
413+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
414+
run: |
415+
rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_FILE_NAME}}.csv | wc -l ) -1))
416+
not_verified=true
417+
row_count=0
418+
while $not_verified; do
419+
row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(DISTINCT(ID)) FROM `${{env.PROJECT_ID}}.${{env.PARQUET_DATASET_ID}}.${{env.INPUT_PARQUET_FILE_NAME}}`')
420+
row_count=$(echo "$row_count_json" | jq -r '.[].f0_')
421+
if [[ "$row_count" == "$rc_orig" ]]; then
422+
echo "PASSED";
423+
not_verified=false;
424+
else
425+
sleep 30s
426+
fi
427+
echo "Got number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count."
428+
done
429+
echo "# records in input Parquet file are: $rc_orig."
430+
echo "Verified number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count."
431+
259432
de-identification-streaming-write:
260433
needs:
261434
- generate-uuid
@@ -282,7 +455,7 @@ jobs:
282455
DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }}
283456
run: |
284457
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
285-
--region=us-central1 \
458+
--region=${{env.REGION}} \
286459
--project=${{env.PROJECT_ID}} \
287460
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
288461
--numWorkers=1 \
@@ -295,7 +468,7 @@ jobs:
295468
--deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \
296469
--batchSize=200000 \
297470
--DLPMethod=DEID \
298-
--serviceAccount=${SERVICE_ACCOUNT_EMAIL} \
471+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
299472
--useStorageWriteApi \
300473
--storageWriteApiTriggeringFrequencySec=2 \
301474
--numStorageWriteApiStreams=2"
@@ -385,7 +558,7 @@ jobs:
385558
DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }}
386559
run: |
387560
gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \
388-
--region=us-central1 \
561+
--region=${{env.REGION}} \
389562
--project=${{env.PROJECT_ID}} \
390563
--tempLocation=gs://${{env.GCS_BUCKET}}/temp \
391564
--numWorkers=1 \
@@ -400,7 +573,7 @@ jobs:
400573
--DLPMethod=REID \
401574
--keyRange=1024 \
402575
--queryPath=gs://${GCS_BUCKET}/${{env.REIDENTIFICATION_QUERY_FILE}} \
403-
--serviceAccount=${SERVICE_ACCOUNT_EMAIL} \
576+
--serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \
404577
--jobName=${{env.REID_JOB_NAME}}"
405578
406579
- name: Verify BQ table
@@ -445,7 +618,9 @@ jobs:
445618
needs:
446619
- generate-uuid
447620
- inspection
621+
- inspect-parquet-data
448622
- de-identification
623+
- deidentify-parquet-data
449624
- de-identification-streaming-write
450625
- re-identification
451626

@@ -459,9 +634,18 @@ jobs:
459634
if: "!cancelled()"
460635
env:
461636
DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }}
637+
PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }}
462638
run: |
463639
bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.DATASET_ID}}
464-
gsutil rm -f gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}_pub_sub.csv
640+
bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.PARQUET_DATASET_ID}}
641+
642+
- name: Clean up pub_sub file
643+
run: |
644+
if ($(gsutil rm -f gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}_pub_sub.csv)); then
645+
echo "Cleared pub_sub file!"
646+
else
647+
echo "pub_sub file not present in storage bucket."
648+
fi
465649
466650
- name: Cancel Inspection Job
467651
if: "!cancelled()"

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ dlp_key.json
2323
!.github/mock-data/*.csv
2424
!.github/mock-data/*.avsc
2525
!.github/mock-data/*.avro
26+
!.github/mock-data/*.parquet
2627

2728
# Package Files #
2829
*.jar

README.md

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -521,6 +521,18 @@ The pipeline supports CSV files with a custom delimiter. The delimiter has to be
521521
gradle run ... -Pargs="... --columnDelimiter=|"
522522
```
523523

524+
#### 5. Parquet
525+
526+
The inspection and de-identification pipelines support Parquet file format where data can be read from GCS storage
527+
bucket and the results of DLP Dataflow pipeline will be written in BigQuery tables. For sample data in Parquet file
528+
format, refer [mock-data](.github/mock-data).
529+
530+
No additional changes are required to run the pipeline except updating the `--filePattern` parameter. For example:
531+
532+
```commandline
533+
gradle run ... -Pargs="... --filePattern=gs://${PROJECT_ID}-demo-data/*.parquet"
534+
```
535+
524536
### Amazon S3 Scanner
525537

526538
To use Amazon S3 as a source of input files, use AWS credentials as instructed below.

build.gradle

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,9 @@ dependencies {
8888
implementation group: 'org.apache.beam', name: 'beam-runners-direct-java', version: dataflowBeamVersion
8989
implementation group: 'org.apache.beam', name: 'beam-sdks-java-extensions-ml', version: dataflowBeamVersion
9090
implementation group: 'org.apache.beam', name: 'beam-sdks-java-io-amazon-web-services', version: dataflowBeamVersion
91+
implementation group: 'org.apache.parquet', name: 'parquet-avro', version: '1.13.1'
92+
implementation group: 'org.apache.parquet', name: 'parquet-hadoop', version: '1.13.1'
93+
implementation group: 'org.apache.beam', name: 'beam-sdks-java-io-parquet', version: '2.14.0'
9194
implementation group: 'org.slf4j', name: 'slf4j-jdk14', version: '2.0.6'
9295
implementation 'com.google.cloud:google-cloud-kms:2.15.0'
9396
implementation 'com.google.guava:guava:31.1-jre'

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
import com.google.swarm.tokenization.common.Util.InputLocation;
3939
import com.google.swarm.tokenization.json.ConvertJsonRecordToDLPRow;
4040
import com.google.swarm.tokenization.json.JsonReaderSplitDoFn;
41+
import com.google.swarm.tokenization.parquet.ParquetReaderSplittableDoFn;
4142
import com.google.swarm.tokenization.txt.ConvertTxtToDLPRow;
4243
import com.google.swarm.tokenization.txt.ParseTextLogDoFn;
4344
import com.google.swarm.tokenization.txt.TxtReaderSplitDoFn;
@@ -216,6 +217,10 @@ private static void runInspectAndDeidPipeline(
216217
ParDo.of(new ConvertTxtToDLPRow(options.getColumnDelimiter(), headers))
217218
.withSideInputs(headers));
218219
break;
220+
case PARQUET:
221+
records = inputFiles
222+
.apply(ParDo.of(new ParquetReaderSplittableDoFn(options.getKeyRange(), options.getSplitSize())));
223+
break;
219224
default:
220225
throw new IllegalArgumentException("Please validate FileType parameter");
221226
}

0 commit comments

Comments
 (0)