From 26fe09d297e5c93c7e1f836318145bc26beb85d3 Mon Sep 17 00:00:00 2001 From: dup05 Date: Fri, 22 Sep 2023 19:58:24 +0530 Subject: [PATCH 1/9] Changes to make queryPath optional --- ...tToBigQueryStreamingV2PipelineOptions.java | 2 +- .../common/BigQueryReadTransform.java | 39 +++++++++++-------- .../swarm/tokenization/common/Util.java | 4 ++ 3 files changed, 28 insertions(+), 17 deletions(-) diff --git a/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2PipelineOptions.java b/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2PipelineOptions.java index a839bf47..4895f468 100644 --- a/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2PipelineOptions.java +++ b/src/main/java/com/google/swarm/tokenization/DLPTextToBigQueryStreamingV2PipelineOptions.java @@ -92,7 +92,7 @@ public interface DLPTextToBigQueryStreamingV2PipelineOptions void setTableRef(String tableRef); @Description("read method default, direct, export") - @Default.Enum("EXPORT") + @Default.Enum("DEFAULT") Method getReadMethod(); void setReadMethod(Method method); diff --git a/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java b/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java index 6676c45e..174bf5be 100644 --- a/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java +++ b/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java @@ -15,6 +15,7 @@ */ package com.google.swarm.tokenization.common; +import autovalue.shaded.org.checkerframework.checker.nullness.qual.Nullable; import com.google.api.services.bigquery.model.TableRow; import com.google.auto.value.AutoValue; import org.apache.beam.sdk.io.gcp.bigquery.BigQueryIO; @@ -50,6 +51,7 @@ public abstract static class Builder { public abstract Builder setKeyRange(Integer keyRange); + @Nullable public abstract Builder setQuery(String query); public abstract BigQueryReadTransform build(); @@ -64,23 +66,28 @@ public PCollection> expand(PBegin input) { switch (readMethod()) { case DEFAULT: - return input - .apply( - "ReadFromBigQuery", - BigQueryIO.readTableRows() - .fromQuery(query()) - .usingStandardSql() - .withMethod(Method.DEFAULT)) - .apply("AddTableNameAsKey", WithKeys.of(tableRef())); + if(query()!=null) { + return input + .apply( + "ReadFromBigQuery", + BigQueryIO.readTableRows() + .fromQuery(query()) + .usingStandardSql() + .withMethod(Method.DEFAULT)) + .apply("AddTableNameAsKey", WithKeys.of(tableRef())); + } + else { + return input + .apply( + "ReadFromBigQuery", + BigQueryIO.readTableRows() + .from(tableRef()) + .withMethod(Method.DEFAULT)) + .apply("AddTableNameAsKey", WithKeys.of(tableRef())); + } + case EXPORT: - return input - .apply( - "ReadFromBigQuery", - BigQueryIO.readTableRows() - .fromQuery(query()) - .usingStandardSql() - .withMethod(Method.DEFAULT)) - .apply("AddTableNameAsKey", WithKeys.of(tableRef())); + throw new IllegalArgumentException("Export method not supported"); case DIRECT_READ: throw new IllegalArgumentException("Direct read not supported"); 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 3d3c4b5e..eca605c1 100644 --- a/src/main/java/com/google/swarm/tokenization/common/Util.java +++ b/src/main/java/com/google/swarm/tokenization/common/Util.java @@ -440,6 +440,10 @@ private static void validateStorageWriteApiAtLeastOnce(BigQueryOptions options) } public static String getQueryFromGcs(String gcsPath) { + if(gcsPath == null){ + LOG.debug("Query path not provided, entire table will be read"); + return null; + } GcsPath path = GcsPath.fromUri(URI.create(gcsPath)); Storage storage = StorageOptions.getDefaultInstance().getService(); BlobId blobId = BlobId.of(path.getBucket(), path.getObject()); From b0bcfa24e526daabf872853d8e03a805851022fd Mon Sep 17 00:00:00 2001 From: dup05 Date: Mon, 9 Oct 2023 11:05:13 +0530 Subject: [PATCH 2/9] Reid changes --- .github/workflows/dlp-pipelines.yml | 88 +++++++++++++++++++++++++++++ 1 file changed, 88 insertions(+) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index b72d867f..b3f06ab0 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -10,6 +10,7 @@ env: PROJECT_ID: "dlp-dataflow-deid-ci-392604" DATASET_ID: "demo_dataset" PARQUET_DATASET_ID: "parquet_results" + REID_DATASET_ID: "reid_results" REGION: "us-central1" GCS_BUCKET: "dlp-dataflow-deid-ci-392604-demo-data" GCS_NOTIFICATION_TOPIC: "projects/dlp-dataflow-deid-ci-392604/topics/dlp-dataflow-deid-ci-gcs-notification-topic" @@ -43,6 +44,7 @@ jobs: output6: ${{ steps.gen-uuid.outputs.parquet_inspect_job_name }} output7: ${{ steps.gen-uuid.outputs.parquet_deid_job_name }} output8: ${{ steps.gen-uuid.outputs.parquet_dataset_id }} + output10: ${{ steps.gen-uuid.outputs.reid_dataset_id }} steps: - name: Generate UUID for workflow @@ -58,6 +60,7 @@ jobs: echo "parquet_inspect_job_name=parquet-inspect-$new_uuid" >> "$GITHUB_OUTPUT" echo "parquet_deid_job_name=parquet-deid-$new_uuid" >> "$GITHUB_OUTPUT" echo "parquet_dataset_id=${{ env.PARQUET_DATASET_ID }}_$modified_uuid" >> "$GITHUB_OUTPUT" + echo "reid_dataset_id=${{ env.REID_DATASET_ID}}_$modified_uuid" >> "$GITHUB_OUTPUT" create-dataset: needs: @@ -73,9 +76,12 @@ jobs: env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | bq --location=US mk -d --description "GitHub CI workflow dataset" ${{ env.DATASET_ID }} bq --location=US mk -d --description "GitHub CI workflow dataset to store parquet results" ${{ env.PARQUET_DATASET_ID }} + bq --location=US mk -d --description "GitHub CI workflow dataset to store reid results" ${{ env.REID_DATASET_ID }} + inspection: needs: @@ -613,6 +619,85 @@ jobs: echo "# records in input query are: $rc_orig." echo "Verified number of rows in ${{env.INPUT_FILE_NAME}}_re_id: $row_count." + re-identification-without-query: + needs: + - generate-uuid + - create-dataset + - de-identification + + runs-on: + - self-hosted + + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v2 + + - name: Setup Java + uses: actions/setup-java@v1 + with: + java-version: 11 + + - name: Setup Gradle + uses: gradle/gradle-build-action@v2 + + - name: Run DLP Pipeline + env: + REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + run: | + gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ + --region=${{env.REGION}} \ + --project=${{env.PROJECT_ID}} \ + --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ + --numWorkers=1 \ + --maxNumWorkers=2 \ + --runner=DataflowRunner \ + --tableRef=${{env.PROJECT_ID}}:${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}} \ + --dataset=${{env.REID_DATASET_ID}} \ + --autoscalingAlgorithm=THROUGHPUT_BASED \ + --workerMachineType=n1-highmem-4 \ + --deidentifyTemplateName=${{env.REID_TEMPLATE_PATH}} \ + --DLPMethod=REID \ + --keyRange=1024 \ + --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}}" + + - name: Verify BQ table + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + not_verified=true + table_count=0 + while $not_verified; do + table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.REID_DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INPUT_FILE_NAME}}_re_id"' | wc -l ) -1)) + if [[ "$table_count" == "1" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + done + echo "Verified number of tables in BQ with id ${{env.INPUT_FILE_NAME}}_re_id: $table_count ." + + - name: Verify distinct rows + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + rc_orig=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}}`') + not_verified=true + row_count=0 + while $not_verified; do + row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.REID_DATASET_ID}}.${{env.INPUT_FILE_NAME}}_re_id`') + row_count=$(echo "$row_count_json" | jq -r '.[].f0_') + if [[ "$row_count" == "$rc_orig" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + done + echo "# records in input query are: $rc_orig." + echo "Verified number of rows in ${{env.INPUT_FILE_NAME}}_re_id: $row_count." + clean-up: if: "!cancelled()" needs: @@ -623,6 +708,7 @@ jobs: - deidentify-parquet-data - de-identification-streaming-write - re-identification + - re-identification-without-query runs-on: - self-hosted @@ -635,9 +721,11 @@ jobs: env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.DATASET_ID}} bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.PARQUET_DATASET_ID}} + bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.REID_DATASET_ID}} - name: Clean up pub_sub file run: | From 211b7825bb346631972b6c6534d74f879ffddb08 Mon Sep 17 00:00:00 2001 From: dup05 Date: Mon, 9 Oct 2023 11:31:29 +0530 Subject: [PATCH 3/9] minor changes --- .github/workflows/dlp-pipelines.yml | 316 +++++++++--------- .../common/BigQueryReadTransform.java | 2 +- 2 files changed, 159 insertions(+), 159 deletions(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index b3f06ab0..2c62994b 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -83,166 +83,166 @@ jobs: bq --location=US mk -d --description "GitHub CI workflow dataset to store reid results" ${{ env.REID_DATASET_ID }} - inspection: - needs: - - generate-uuid - - create-dataset - - runs-on: - - self-hosted - - timeout-minutes: 30 - - steps: - - uses: actions/checkout@v2 - - - name: Setup Java - uses: actions/setup-java@v1 - with: - java-version: 11 - - - name: Setup Gradle - uses: gradle/gradle-build-action@v2 - - - name: Run DLP Pipeline - env: - INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output2 }} - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ - --streaming --enableStreamingEngine \ - --region=${{env.REGION}} \ - --project=${{env.PROJECT_ID}} \ - --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ - --numWorkers=1 \ - --maxNumWorkers=2 \ - --runner=DataflowRunner \ - --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}*.csv \ - --dataset=${{env.DATASET_ID}} \ - --workerMachineType=n1-highmem-4 \ - --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ - --batchSize=200000 \ - --DLPMethod=INSPECT \ - --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ - --jobName=${{env.INSPECT_JOB_NAME}} \ - --gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}" - sleep 30s - gsutil cp gs://${{env.GCS_BUCKET}}/temp/csv/tiny_csv_pub_sub.csv gs://${{env.GCS_BUCKET}} - - - name: Verify BQ table - env: - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - not_verified=true - table_count=0 - while $not_verified; do - table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INSPECTION_TABLE_ID}}"' | wc -l ) -1)) - if [[ "$table_count" == "1" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." - done - echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." - - - name: Verify distinct rows - env: - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - not_verified=true - row_count=0 - while $not_verified; do - row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(*) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INSPECTION_TABLE_ID}}`') - row_count=$(echo "$row_count_json" | jq -r '.[].f0_') - if [[ "$row_count" -gt ${{env.NUM_INSPECTION_RECORDS_THRESHOLD}} ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." - done - echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." +# inspection: +# needs: +# - generate-uuid +# - create-dataset +# +# runs-on: +# - self-hosted +# +# timeout-minutes: 30 +# +# steps: +# - uses: actions/checkout@v2 +# +# - name: Setup Java +# uses: actions/setup-java@v1 +# with: +# java-version: 11 +# +# - name: Setup Gradle +# uses: gradle/gradle-build-action@v2 +# +# - name: Run DLP Pipeline +# env: +# INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output2 }} +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ +# --streaming --enableStreamingEngine \ +# --region=${{env.REGION}} \ +# --project=${{env.PROJECT_ID}} \ +# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ +# --numWorkers=1 \ +# --maxNumWorkers=2 \ +# --runner=DataflowRunner \ +# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}*.csv \ +# --dataset=${{env.DATASET_ID}} \ +# --workerMachineType=n1-highmem-4 \ +# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ +# --batchSize=200000 \ +# --DLPMethod=INSPECT \ +# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ +# --jobName=${{env.INSPECT_JOB_NAME}} \ +# --gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}" +# sleep 30s +# gsutil cp gs://${{env.GCS_BUCKET}}/temp/csv/tiny_csv_pub_sub.csv gs://${{env.GCS_BUCKET}} +# +# - name: Verify BQ table +# env: +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# not_verified=true +# table_count=0 +# while $not_verified; do +# table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INSPECTION_TABLE_ID}}"' | wc -l ) -1)) +# if [[ "$table_count" == "1" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." +# done +# echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." +# +# - name: Verify distinct rows +# env: +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# not_verified=true +# row_count=0 +# while $not_verified; do +# row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(*) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INSPECTION_TABLE_ID}}`') +# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') +# if [[ "$row_count" -gt ${{env.NUM_INSPECTION_RECORDS_THRESHOLD}} ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." +# done +# echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." # Inspect only existing parquet files - inspect-parquet-data: - needs: - - generate-uuid - - create-dataset - - runs-on: - - self-hosted - - timeout-minutes: 30 - - steps: - - uses: actions/checkout@v2 - - - name: Setup Java - uses: actions/setup-java@v1 - with: - java-version: 11 - - - name: Setup Gradle - uses: gradle/gradle-build-action@v2 - - - name: Run DLP Pipeline - env: - PARQUET_INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output6 }} - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ - --region=${{env.REGION}} \ - --project=${{env.PROJECT_ID}} \ - --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ - --numWorkers=1 --maxNumWorkers=2 \ - --runner=DataflowRunner \ - --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}*.parquet \ - --dataset=${{env.PARQUET_DATASET_ID}} \ - --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ - --batchSize=200000 \ - --DLPMethod=INSPECT \ - --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ - --jobName=${{env.PARQUET_INSPECT_JOB_NAME}}" - - - name: Verify BQ table - env: - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - not_verified=true - table_count=0 - while $not_verified; do - 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)) - if [[ "$table_count" == "1" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." - done - echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." - - - name: Verify distinct rows - env: - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - not_verified=true - row_count=0 - while $not_verified; do - 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}}`') - row_count=$(echo "$row_count_json" | jq -r '.[].f0_') - if [[ "$row_count" == ${{env.PARQUET_INSPECTION_RECORDS_THRESHOLD}} ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." - done - echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." +# inspect-parquet-data: +# needs: +# - generate-uuid +# - create-dataset +# +# runs-on: +# - self-hosted +# +# timeout-minutes: 30 +# +# steps: +# - uses: actions/checkout@v2 +# +# - name: Setup Java +# uses: actions/setup-java@v1 +# with: +# java-version: 11 +# +# - name: Setup Gradle +# uses: gradle/gradle-build-action@v2 +# +# - name: Run DLP Pipeline +# env: +# PARQUET_INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output6 }} +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ +# --region=${{env.REGION}} \ +# --project=${{env.PROJECT_ID}} \ +# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ +# --numWorkers=1 --maxNumWorkers=2 \ +# --runner=DataflowRunner \ +# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}*.parquet \ +# --dataset=${{env.PARQUET_DATASET_ID}} \ +# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ +# --batchSize=200000 \ +# --DLPMethod=INSPECT \ +# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ +# --jobName=${{env.PARQUET_INSPECT_JOB_NAME}}" +# +# - name: Verify BQ table +# env: +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# not_verified=true +# table_count=0 +# while $not_verified; do +# 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)) +# if [[ "$table_count" == "1" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." +# done +# echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." +# +# - name: Verify distinct rows +# env: +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# not_verified=true +# row_count=0 +# while $not_verified; do +# 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}}`') +# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') +# if [[ "$row_count" == ${{env.PARQUET_INSPECTION_RECORDS_THRESHOLD}} ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." +# done +# echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." de-identification: needs: diff --git a/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java b/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java index 174bf5be..9c066e5d 100644 --- a/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java +++ b/src/main/java/com/google/swarm/tokenization/common/BigQueryReadTransform.java @@ -40,6 +40,7 @@ public abstract class BigQueryReadTransform public abstract Integer keyRange(); + @Nullable public abstract String query(); @AutoValue.Builder @@ -51,7 +52,6 @@ public abstract static class Builder { public abstract Builder setKeyRange(Integer keyRange); - @Nullable public abstract Builder setQuery(String query); public abstract BigQueryReadTransform build(); From 6c31dfc41747f657bc33c0f02abd5b2bc83ee457 Mon Sep 17 00:00:00 2001 From: dup05 Date: Mon, 9 Oct 2023 11:32:43 +0530 Subject: [PATCH 4/9] temp changes to test CI --- .github/workflows/dlp-pipelines.yml | 330 ++++++++++++++-------------- 1 file changed, 165 insertions(+), 165 deletions(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index 2c62994b..b5834194 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -355,168 +355,168 @@ jobs: echo "Verified number of rows in ${{env.INPUT_FILE_NAME}}_pub_sub: $row_count." # Deidentify only existing parquet files - deidentify-parquet-data: - needs: - - generate-uuid - - create-dataset - - runs-on: - - self-hosted - - timeout-minutes: 30 - - steps: - - uses: actions/checkout@v2 - - - name: Setup Java - uses: actions/setup-java@v1 - with: - java-version: 11 - - - name: Setup Gradle - uses: gradle/gradle-build-action@v2 - - - name: Run DLP Pipeline - env: - PARQUET_DEID_JOB_NAME: ${{ needs.generate-uuid.outputs.output7 }} - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ - --region=${{env.REGION}} \ - --project=${{env.PROJECT_ID}} \ - --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ - --numWorkers=1 --maxNumWorkers=2 \ - --runner=DataflowRunner \ - --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}.parquet \ - --dataset=${{env.PARQUET_DATASET_ID}} \ - --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ - --deidentifyTemplateName=${{env.PARQUET_DEID_TEMPLATE_PATH}} \ - --batchSize=200000 \ - --DLPMethod=DEID \ - --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ - --jobName=${{env.PARQUET_DEID_JOB_NAME}}" - - - name: Verify BQ tables - env: - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - not_verified=true - table_count=0 - while $not_verified; do - 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)) - if [[ "$table_count" == "1" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." - done - echo "Verified number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." - - - name: Verify distinct rows of existing file - env: - PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - run: | - rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_FILE_NAME}}.csv | wc -l ) -1)) - not_verified=true - row_count=0 - while $not_verified; do - 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}}`') - row_count=$(echo "$row_count_json" | jq -r '.[].f0_') - if [[ "$row_count" == "$rc_orig" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." - done - echo "# records in input Parquet file are: $rc_orig." - echo "Verified number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." - - de-identification-streaming-write: - needs: - - generate-uuid - - create-dataset - - runs-on: - - self-hosted - - timeout-minutes: 30 - - steps: - - uses: actions/checkout@v2 - - - name: Setup Java - uses: actions/setup-java@v1 - with: - java-version: 11 - - - name: Setup Gradle - uses: gradle/gradle-build-action@v2 - - - name: Run DLP Pipeline for streaming write - env: - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ - --region=${{env.REGION}} \ - --project=${{env.PROJECT_ID}} \ - --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ - --numWorkers=1 \ - --maxNumWorkers=2 \ - --runner=DataflowRunner \ - --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv \ - --dataset=${{env.DATASET_ID}} \ - --workerMachineType=n1-highmem-4 \ - --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ - --deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \ - --batchSize=200000 \ - --DLPMethod=DEID \ - --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ - --useStorageWriteApi \ - --storageWriteApiTriggeringFrequencySec=2 \ - --numStorageWriteApiStreams=2" - - - name: Verify BQ table for streaming write - env: - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - not_verified=true - table_count=0 - while $not_verified; do - table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id LIKE "${{env.INPUT_STREAMING_WRITE_FILE_NAME}}%"' | wc -l ) -1)) - if [[ "$table_count" == "1" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." - done - echo "Verified number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." +# deidentify-parquet-data: +# needs: +# - generate-uuid +# - create-dataset +# +# runs-on: +# - self-hosted +# +# timeout-minutes: 30 +# +# steps: +# - uses: actions/checkout@v2 +# +# - name: Setup Java +# uses: actions/setup-java@v1 +# with: +# java-version: 11 +# +# - name: Setup Gradle +# uses: gradle/gradle-build-action@v2 +# +# - name: Run DLP Pipeline +# env: +# PARQUET_DEID_JOB_NAME: ${{ needs.generate-uuid.outputs.output7 }} +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ +# --region=${{env.REGION}} \ +# --project=${{env.PROJECT_ID}} \ +# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ +# --numWorkers=1 --maxNumWorkers=2 \ +# --runner=DataflowRunner \ +# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}.parquet \ +# --dataset=${{env.PARQUET_DATASET_ID}} \ +# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ +# --deidentifyTemplateName=${{env.PARQUET_DEID_TEMPLATE_PATH}} \ +# --batchSize=200000 \ +# --DLPMethod=DEID \ +# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ +# --jobName=${{env.PARQUET_DEID_JOB_NAME}}" +# +# - name: Verify BQ tables +# env: +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# not_verified=true +# table_count=0 +# while $not_verified; do +# 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)) +# if [[ "$table_count" == "1" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." +# done +# echo "Verified number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." +# +# - name: Verify distinct rows of existing file +# env: +# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} +# run: | +# rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_FILE_NAME}}.csv | wc -l ) -1)) +# not_verified=true +# row_count=0 +# while $not_verified; do +# 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}}`') +# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') +# if [[ "$row_count" == "$rc_orig" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." +# done +# echo "# records in input Parquet file are: $rc_orig." +# echo "Verified number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." - - name: Verify distinct rows of streaming write file - env: - DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - run: | - rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv | wc -l ) -1)) - not_verified=true - row_count=0 - while $not_verified; do - row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(DISTINCT(ID)) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_STREAMING_WRITE_FILE_NAME}}`') - row_count=$(echo "$row_count_json" | jq -r '.[].f0_') - if [[ "$row_count" == "$rc_orig" ]]; then - echo "PASSED"; - not_verified=false; - else - sleep 30s - fi - echo "Got number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." - done - echo "# records in input CSV file are: $rc_orig." - echo "Verified number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." +# de-identification-streaming-write: +# needs: +# - generate-uuid +# - create-dataset +# +# runs-on: +# - self-hosted +# +# timeout-minutes: 30 +# +# steps: +# - uses: actions/checkout@v2 +# +# - name: Setup Java +# uses: actions/setup-java@v1 +# with: +# java-version: 11 +# +# - name: Setup Gradle +# uses: gradle/gradle-build-action@v2 +# +# - name: Run DLP Pipeline for streaming write +# env: +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ +# --region=${{env.REGION}} \ +# --project=${{env.PROJECT_ID}} \ +# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ +# --numWorkers=1 \ +# --maxNumWorkers=2 \ +# --runner=DataflowRunner \ +# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv \ +# --dataset=${{env.DATASET_ID}} \ +# --workerMachineType=n1-highmem-4 \ +# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ +# --deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \ +# --batchSize=200000 \ +# --DLPMethod=DEID \ +# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ +# --useStorageWriteApi \ +# --storageWriteApiTriggeringFrequencySec=2 \ +# --numStorageWriteApiStreams=2" +# +# - name: Verify BQ table for streaming write +# env: +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# not_verified=true +# table_count=0 +# while $not_verified; do +# table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id LIKE "${{env.INPUT_STREAMING_WRITE_FILE_NAME}}%"' | wc -l ) -1)) +# if [[ "$table_count" == "1" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." +# done +# echo "Verified number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." +# +# - name: Verify distinct rows of streaming write file +# env: +# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} +# run: | +# rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv | wc -l ) -1)) +# not_verified=true +# row_count=0 +# while $not_verified; do +# row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(DISTINCT(ID)) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_STREAMING_WRITE_FILE_NAME}}`') +# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') +# if [[ "$row_count" == "$rc_orig" ]]; then +# echo "PASSED"; +# not_verified=false; +# else +# sleep 30s +# fi +# echo "Got number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." +# done +# echo "# records in input CSV file are: $rc_orig." +# echo "Verified number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." re-identification: needs: @@ -702,11 +702,11 @@ jobs: if: "!cancelled()" needs: - generate-uuid - - inspection - - inspect-parquet-data +# - inspection +# - inspect-parquet-data - de-identification - - deidentify-parquet-data - - de-identification-streaming-write +# - deidentify-parquet-data +# - de-identification-streaming-write - re-identification - re-identification-without-query From 1ffc3b9e10d539624b2f2a47cebc1959d9e41f21 Mon Sep 17 00:00:00 2001 From: dup05 Date: Mon, 9 Oct 2023 11:59:58 +0530 Subject: [PATCH 5/9] CI changes --- .github/workflows/dlp-pipelines.yml | 3 +++ 1 file changed, 3 insertions(+) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index b5834194..5d64040d 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -643,6 +643,7 @@ jobs: - name: Run DLP Pipeline env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ @@ -664,6 +665,7 @@ jobs: - name: Verify BQ table env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | not_verified=true table_count=0 @@ -681,6 +683,7 @@ jobs: - name: Verify distinct rows env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | rc_orig=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}}`') not_verified=true From 40330a4ed3c590ef510552b54e2b224d0e77bf15 Mon Sep 17 00:00:00 2001 From: dup05 Date: Wed, 11 Oct 2023 10:33:32 +0530 Subject: [PATCH 6/9] Fixes in CI workflow --- .github/workflows/dlp-pipelines.yml | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index 5d64040d..e46aad2f 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -685,7 +685,8 @@ jobs: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | - rc_orig=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}}`') + rc_orig_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}}`') + rc_orig=$(echo "$rc_orig_json" | jq -r '.[].f0_') not_verified=true row_count=0 while $not_verified; do From 21b45f4ab3dd305d76f053192581288a9fbc80d6 Mon Sep 17 00:00:00 2001 From: dup05 Date: Wed, 11 Oct 2023 10:52:48 +0530 Subject: [PATCH 7/9] removed temp changes --- .github/workflows/dlp-pipelines.yml | 640 ++++++++++++++-------------- 1 file changed, 320 insertions(+), 320 deletions(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index e46aad2f..837973b0 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -83,166 +83,166 @@ jobs: bq --location=US mk -d --description "GitHub CI workflow dataset to store reid results" ${{ env.REID_DATASET_ID }} -# inspection: -# needs: -# - generate-uuid -# - create-dataset -# -# runs-on: -# - self-hosted -# -# timeout-minutes: 30 -# -# steps: -# - uses: actions/checkout@v2 -# -# - name: Setup Java -# uses: actions/setup-java@v1 -# with: -# java-version: 11 -# -# - name: Setup Gradle -# uses: gradle/gradle-build-action@v2 -# -# - name: Run DLP Pipeline -# env: -# INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output2 }} -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ -# --streaming --enableStreamingEngine \ -# --region=${{env.REGION}} \ -# --project=${{env.PROJECT_ID}} \ -# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ -# --numWorkers=1 \ -# --maxNumWorkers=2 \ -# --runner=DataflowRunner \ -# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}*.csv \ -# --dataset=${{env.DATASET_ID}} \ -# --workerMachineType=n1-highmem-4 \ -# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ -# --batchSize=200000 \ -# --DLPMethod=INSPECT \ -# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ -# --jobName=${{env.INSPECT_JOB_NAME}} \ -# --gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}" -# sleep 30s -# gsutil cp gs://${{env.GCS_BUCKET}}/temp/csv/tiny_csv_pub_sub.csv gs://${{env.GCS_BUCKET}} -# -# - name: Verify BQ table -# env: -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# not_verified=true -# table_count=0 -# while $not_verified; do -# table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INSPECTION_TABLE_ID}}"' | wc -l ) -1)) -# if [[ "$table_count" == "1" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." -# done -# echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." -# -# - name: Verify distinct rows -# env: -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# not_verified=true -# row_count=0 -# while $not_verified; do -# row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(*) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INSPECTION_TABLE_ID}}`') -# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') -# if [[ "$row_count" -gt ${{env.NUM_INSPECTION_RECORDS_THRESHOLD}} ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." -# done -# echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." + inspection: + needs: + - generate-uuid + - create-dataset + + runs-on: + - self-hosted + + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v2 + + - name: Setup Java + uses: actions/setup-java@v1 + with: + java-version: 11 + + - name: Setup Gradle + uses: gradle/gradle-build-action@v2 + + - name: Run DLP Pipeline + env: + INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output2 }} + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ + --streaming --enableStreamingEngine \ + --region=${{env.REGION}} \ + --project=${{env.PROJECT_ID}} \ + --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ + --numWorkers=1 \ + --maxNumWorkers=2 \ + --runner=DataflowRunner \ + --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_FILE_NAME}}*.csv \ + --dataset=${{env.DATASET_ID}} \ + --workerMachineType=n1-highmem-4 \ + --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ + --batchSize=200000 \ + --DLPMethod=INSPECT \ + --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ + --jobName=${{env.INSPECT_JOB_NAME}} \ + --gcsNotificationTopic=${{env.GCS_NOTIFICATION_TOPIC}}" + sleep 30s + gsutil cp gs://${{env.GCS_BUCKET}}/temp/csv/tiny_csv_pub_sub.csv gs://${{env.GCS_BUCKET}} + + - name: Verify BQ table + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + not_verified=true + table_count=0 + while $not_verified; do + table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INSPECTION_TABLE_ID}}"' | wc -l ) -1)) + if [[ "$table_count" == "1" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." + done + echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." + + - name: Verify distinct rows + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + not_verified=true + row_count=0 + while $not_verified; do + row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(*) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INSPECTION_TABLE_ID}}`') + row_count=$(echo "$row_count_json" | jq -r '.[].f0_') + if [[ "$row_count" -gt ${{env.NUM_INSPECTION_RECORDS_THRESHOLD}} ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." + done + echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." # Inspect only existing parquet files -# inspect-parquet-data: -# needs: -# - generate-uuid -# - create-dataset -# -# runs-on: -# - self-hosted -# -# timeout-minutes: 30 -# -# steps: -# - uses: actions/checkout@v2 -# -# - name: Setup Java -# uses: actions/setup-java@v1 -# with: -# java-version: 11 -# -# - name: Setup Gradle -# uses: gradle/gradle-build-action@v2 -# -# - name: Run DLP Pipeline -# env: -# PARQUET_INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output6 }} -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ -# --region=${{env.REGION}} \ -# --project=${{env.PROJECT_ID}} \ -# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ -# --numWorkers=1 --maxNumWorkers=2 \ -# --runner=DataflowRunner \ -# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}*.parquet \ -# --dataset=${{env.PARQUET_DATASET_ID}} \ -# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ -# --batchSize=200000 \ -# --DLPMethod=INSPECT \ -# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ -# --jobName=${{env.PARQUET_INSPECT_JOB_NAME}}" -# -# - name: Verify BQ table -# env: -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# not_verified=true -# table_count=0 -# while $not_verified; do -# 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)) -# if [[ "$table_count" == "1" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." -# done -# echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." -# -# - name: Verify distinct rows -# env: -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# not_verified=true -# row_count=0 -# while $not_verified; do -# 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}}`') -# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') -# if [[ "$row_count" == ${{env.PARQUET_INSPECTION_RECORDS_THRESHOLD}} ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." -# done -# echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." + inspect-parquet-data: + needs: + - generate-uuid + - create-dataset + + runs-on: + - self-hosted + + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v2 + + - name: Setup Java + uses: actions/setup-java@v1 + with: + java-version: 11 + + - name: Setup Gradle + uses: gradle/gradle-build-action@v2 + + - name: Run DLP Pipeline + env: + PARQUET_INSPECT_JOB_NAME: ${{ needs.generate-uuid.outputs.output6 }} + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ + --region=${{env.REGION}} \ + --project=${{env.PROJECT_ID}} \ + --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ + --numWorkers=1 --maxNumWorkers=2 \ + --runner=DataflowRunner \ + --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}*.parquet \ + --dataset=${{env.PARQUET_DATASET_ID}} \ + --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ + --batchSize=200000 \ + --DLPMethod=INSPECT \ + --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ + --jobName=${{env.PARQUET_INSPECT_JOB_NAME}}" + + - name: Verify BQ table + env: + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + not_verified=true + table_count=0 + while $not_verified; do + 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)) + if [[ "$table_count" == "1" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." + done + echo "Verified number of tables in BQ with id ${{env.INSPECTION_TABLE_ID}}: $table_count ." + + - name: Verify distinct rows + env: + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + not_verified=true + row_count=0 + while $not_verified; do + 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}}`') + row_count=$(echo "$row_count_json" | jq -r '.[].f0_') + if [[ "$row_count" == ${{env.PARQUET_INSPECTION_RECORDS_THRESHOLD}} ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." + done + echo "Verified number of rows in ${{env.INSPECTION_TABLE_ID}}: $row_count." de-identification: needs: @@ -355,168 +355,168 @@ jobs: echo "Verified number of rows in ${{env.INPUT_FILE_NAME}}_pub_sub: $row_count." # Deidentify only existing parquet files -# deidentify-parquet-data: -# needs: -# - generate-uuid -# - create-dataset -# -# runs-on: -# - self-hosted -# -# timeout-minutes: 30 -# -# steps: -# - uses: actions/checkout@v2 -# -# - name: Setup Java -# uses: actions/setup-java@v1 -# with: -# java-version: 11 -# -# - name: Setup Gradle -# uses: gradle/gradle-build-action@v2 -# -# - name: Run DLP Pipeline -# env: -# PARQUET_DEID_JOB_NAME: ${{ needs.generate-uuid.outputs.output7 }} -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ -# --region=${{env.REGION}} \ -# --project=${{env.PROJECT_ID}} \ -# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ -# --numWorkers=1 --maxNumWorkers=2 \ -# --runner=DataflowRunner \ -# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}.parquet \ -# --dataset=${{env.PARQUET_DATASET_ID}} \ -# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ -# --deidentifyTemplateName=${{env.PARQUET_DEID_TEMPLATE_PATH}} \ -# --batchSize=200000 \ -# --DLPMethod=DEID \ -# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ -# --jobName=${{env.PARQUET_DEID_JOB_NAME}}" -# -# - name: Verify BQ tables -# env: -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# not_verified=true -# table_count=0 -# while $not_verified; do -# 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)) -# if [[ "$table_count" == "1" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." -# done -# echo "Verified number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." -# -# - name: Verify distinct rows of existing file -# env: -# PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} -# run: | -# rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_FILE_NAME}}.csv | wc -l ) -1)) -# not_verified=true -# row_count=0 -# while $not_verified; do -# 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}}`') -# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') -# if [[ "$row_count" == "$rc_orig" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." -# done -# echo "# records in input Parquet file are: $rc_orig." -# echo "Verified number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." - -# de-identification-streaming-write: -# needs: -# - generate-uuid -# - create-dataset -# -# runs-on: -# - self-hosted -# -# timeout-minutes: 30 -# -# steps: -# - uses: actions/checkout@v2 -# -# - name: Setup Java -# uses: actions/setup-java@v1 -# with: -# java-version: 11 -# -# - name: Setup Gradle -# uses: gradle/gradle-build-action@v2 -# -# - name: Run DLP Pipeline for streaming write -# env: -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ -# --region=${{env.REGION}} \ -# --project=${{env.PROJECT_ID}} \ -# --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ -# --numWorkers=1 \ -# --maxNumWorkers=2 \ -# --runner=DataflowRunner \ -# --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv \ -# --dataset=${{env.DATASET_ID}} \ -# --workerMachineType=n1-highmem-4 \ -# --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ -# --deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \ -# --batchSize=200000 \ -# --DLPMethod=DEID \ -# --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ -# --useStorageWriteApi \ -# --storageWriteApiTriggeringFrequencySec=2 \ -# --numStorageWriteApiStreams=2" -# -# - name: Verify BQ table for streaming write -# env: -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# not_verified=true -# table_count=0 -# while $not_verified; do -# table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id LIKE "${{env.INPUT_STREAMING_WRITE_FILE_NAME}}%"' | wc -l ) -1)) -# if [[ "$table_count" == "1" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." -# done -# echo "Verified number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." -# -# - name: Verify distinct rows of streaming write file -# env: -# DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} -# run: | -# rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv | wc -l ) -1)) -# not_verified=true -# row_count=0 -# while $not_verified; do -# row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(DISTINCT(ID)) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_STREAMING_WRITE_FILE_NAME}}`') -# row_count=$(echo "$row_count_json" | jq -r '.[].f0_') -# if [[ "$row_count" == "$rc_orig" ]]; then -# echo "PASSED"; -# not_verified=false; -# else -# sleep 30s -# fi -# echo "Got number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." -# done -# echo "# records in input CSV file are: $rc_orig." -# echo "Verified number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." + deidentify-parquet-data: + needs: + - generate-uuid + - create-dataset + + runs-on: + - self-hosted + + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v2 + + - name: Setup Java + uses: actions/setup-java@v1 + with: + java-version: 11 + + - name: Setup Gradle + uses: gradle/gradle-build-action@v2 + + - name: Run DLP Pipeline + env: + PARQUET_DEID_JOB_NAME: ${{ needs.generate-uuid.outputs.output7 }} + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ + --region=${{env.REGION}} \ + --project=${{env.PROJECT_ID}} \ + --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ + --numWorkers=1 --maxNumWorkers=2 \ + --runner=DataflowRunner \ + --filePattern=gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_PARQUET_FILE_NAME}}.parquet \ + --dataset=${{env.PARQUET_DATASET_ID}} \ + --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ + --deidentifyTemplateName=${{env.PARQUET_DEID_TEMPLATE_PATH}} \ + --batchSize=200000 \ + --DLPMethod=DEID \ + --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ + --jobName=${{env.PARQUET_DEID_JOB_NAME}}" + + - name: Verify BQ tables + env: + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + not_verified=true + table_count=0 + while $not_verified; do + 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)) + if [[ "$table_count" == "1" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." + done + echo "Verified number of tables in BQ with id ${{env.INPUT_PARQUET_FILE_NAME}}*: $table_count ." + + - name: Verify distinct rows of existing file + env: + PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} + run: | + rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.SAMPLE_DATA_DIR}}/${{env.INPUT_FILE_NAME}}.csv | wc -l ) -1)) + not_verified=true + row_count=0 + while $not_verified; do + 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}}`') + row_count=$(echo "$row_count_json" | jq -r '.[].f0_') + if [[ "$row_count" == "$rc_orig" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." + done + echo "# records in input Parquet file are: $rc_orig." + echo "Verified number of rows in ${{env.INPUT_PARQUET_FILE_NAME}}: $row_count." + + de-identification-streaming-write: + needs: + - generate-uuid + - create-dataset + + runs-on: + - self-hosted + + timeout-minutes: 30 + + steps: + - uses: actions/checkout@v2 + + - name: Setup Java + uses: actions/setup-java@v1 + with: + java-version: 11 + + - name: Setup Gradle + uses: gradle/gradle-build-action@v2 + + - name: Run DLP Pipeline for streaming write + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ + --region=${{env.REGION}} \ + --project=${{env.PROJECT_ID}} \ + --tempLocation=gs://${{env.GCS_BUCKET}}/temp \ + --numWorkers=1 \ + --maxNumWorkers=2 \ + --runner=DataflowRunner \ + --filePattern=gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv \ + --dataset=${{env.DATASET_ID}} \ + --workerMachineType=n1-highmem-4 \ + --inspectTemplateName=${{env.INSPECT_TEMPLATE_PATH}} \ + --deidentifyTemplateName=${{env.DEID_TEMPLATE_PATH}} \ + --batchSize=200000 \ + --DLPMethod=DEID \ + --serviceAccount=${{env.SERVICE_ACCOUNT_EMAIL}} \ + --useStorageWriteApi \ + --storageWriteApiTriggeringFrequencySec=2 \ + --numStorageWriteApiStreams=2" + + - name: Verify BQ table for streaming write + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + not_verified=true + table_count=0 + while $not_verified; do + table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}`.__TABLES__ WHERE table_id LIKE "${{env.INPUT_STREAMING_WRITE_FILE_NAME}}%"' | wc -l ) -1)) + if [[ "$table_count" == "1" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." + done + echo "Verified number of tables in BQ with id ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $table_count ." + + - name: Verify distinct rows of streaming write file + env: + DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} + run: | + rc_orig=$(($(gcloud storage cat gs://${{env.GCS_BUCKET}}/${{env.INPUT_STREAMING_WRITE_FILE_NAME}}.csv | wc -l ) -1)) + not_verified=true + row_count=0 + while $not_verified; do + row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(DISTINCT(ID)) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_STREAMING_WRITE_FILE_NAME}}`') + row_count=$(echo "$row_count_json" | jq -r '.[].f0_') + if [[ "$row_count" == "$rc_orig" ]]; then + echo "PASSED"; + not_verified=false; + else + sleep 30s + fi + echo "Got number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." + done + echo "# records in input CSV file are: $rc_orig." + echo "Verified number of rows in ${{env.INPUT_STREAMING_WRITE_FILE_NAME}}: $row_count." re-identification: needs: From 855e5e1a4bb426bbe84136201a68c5df6b5fa97b Mon Sep 17 00:00:00 2001 From: dup05 Date: Wed, 11 Oct 2023 10:55:06 +0530 Subject: [PATCH 8/9] removed temp changes --- .github/workflows/dlp-pipelines.yml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index 837973b0..15c5e818 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -706,11 +706,11 @@ jobs: if: "!cancelled()" needs: - generate-uuid -# - inspection -# - inspect-parquet-data + - inspection + - inspect-parquet-data - de-identification -# - deidentify-parquet-data -# - de-identification-streaming-write + - deidentify-parquet-data + - de-identification-streaming-write - re-identification - re-identification-without-query From 638f2ba40b08c85ad7a9d1b7c728b5ce0ed37dea Mon Sep 17 00:00:00 2001 From: dup05 Date: Tue, 17 Oct 2023 11:08:42 +0530 Subject: [PATCH 9/9] readme changes --- .github/workflows/dlp-pipelines.yml | 26 +++++++++++++------------- README.md | 4 +++- 2 files changed, 16 insertions(+), 14 deletions(-) diff --git a/.github/workflows/dlp-pipelines.yml b/.github/workflows/dlp-pipelines.yml index 15c5e818..a013e497 100644 --- a/.github/workflows/dlp-pipelines.yml +++ b/.github/workflows/dlp-pipelines.yml @@ -10,7 +10,7 @@ env: PROJECT_ID: "dlp-dataflow-deid-ci-392604" DATASET_ID: "demo_dataset" PARQUET_DATASET_ID: "parquet_results" - REID_DATASET_ID: "reid_results" + REID_WITHOUT_QUERY_DATASET_ID: "reid_without_query_results" REGION: "us-central1" GCS_BUCKET: "dlp-dataflow-deid-ci-392604-demo-data" GCS_NOTIFICATION_TOPIC: "projects/dlp-dataflow-deid-ci-392604/topics/dlp-dataflow-deid-ci-gcs-notification-topic" @@ -44,7 +44,7 @@ jobs: output6: ${{ steps.gen-uuid.outputs.parquet_inspect_job_name }} output7: ${{ steps.gen-uuid.outputs.parquet_deid_job_name }} output8: ${{ steps.gen-uuid.outputs.parquet_dataset_id }} - output10: ${{ steps.gen-uuid.outputs.reid_dataset_id }} + output10: ${{ steps.gen-uuid.outputs.reid_without_query_dataset_id }} steps: - name: Generate UUID for workflow @@ -60,7 +60,7 @@ jobs: echo "parquet_inspect_job_name=parquet-inspect-$new_uuid" >> "$GITHUB_OUTPUT" echo "parquet_deid_job_name=parquet-deid-$new_uuid" >> "$GITHUB_OUTPUT" echo "parquet_dataset_id=${{ env.PARQUET_DATASET_ID }}_$modified_uuid" >> "$GITHUB_OUTPUT" - echo "reid_dataset_id=${{ env.REID_DATASET_ID}}_$modified_uuid" >> "$GITHUB_OUTPUT" + echo "reid_without_query_dataset_id=${{ env.REID_WITHOUT_QUERY_DATASET_ID}}_$modified_uuid" >> "$GITHUB_OUTPUT" create-dataset: needs: @@ -76,11 +76,11 @@ jobs: env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + REID_WITHOUT_QUERY_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | bq --location=US mk -d --description "GitHub CI workflow dataset" ${{ env.DATASET_ID }} bq --location=US mk -d --description "GitHub CI workflow dataset to store parquet results" ${{ env.PARQUET_DATASET_ID }} - bq --location=US mk -d --description "GitHub CI workflow dataset to store reid results" ${{ env.REID_DATASET_ID }} + bq --location=US mk -d --description "GitHub CI workflow dataset to store reid results" ${{ env.REID_WITHOUT_QUERY_DATASET_ID }} inspection: @@ -644,7 +644,7 @@ jobs: - name: Run DLP Pipeline env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + REID_WITHOUT_QUERY_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | gradle run -DmainClass=com.google.swarm.tokenization.DLPTextToBigQueryStreamingV2 -Pargs=" \ --region=${{env.REGION}} \ @@ -654,7 +654,7 @@ jobs: --maxNumWorkers=2 \ --runner=DataflowRunner \ --tableRef=${{env.PROJECT_ID}}:${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}} \ - --dataset=${{env.REID_DATASET_ID}} \ + --dataset=${{env.REID_WITHOUT_QUERY_DATASET_ID}} \ --autoscalingAlgorithm=THROUGHPUT_BASED \ --workerMachineType=n1-highmem-4 \ --deidentifyTemplateName=${{env.REID_TEMPLATE_PATH}} \ @@ -665,12 +665,12 @@ jobs: - name: Verify BQ table env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + REID_WITHOUT_QUERY_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | not_verified=true table_count=0 while $not_verified; do - table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.REID_DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INPUT_FILE_NAME}}_re_id"' | wc -l ) -1)) + table_count=$(($(bq query --use_legacy_sql=false --format csv 'SELECT * FROM `${{env.PROJECT_ID}}.${{env.REID_WITHOUT_QUERY_DATASET_ID}}`.__TABLES__ WHERE table_id="${{env.INPUT_FILE_NAME}}_re_id"' | wc -l ) -1)) if [[ "$table_count" == "1" ]]; then echo "PASSED"; not_verified=false; @@ -683,14 +683,14 @@ jobs: - name: Verify distinct rows env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} - REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + REID_WITHOUT_QUERY_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | rc_orig_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.DATASET_ID}}.${{env.INPUT_FILE_NAME}}`') rc_orig=$(echo "$rc_orig_json" | jq -r '.[].f0_') not_verified=true row_count=0 while $not_verified; do - row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.REID_DATASET_ID}}.${{env.INPUT_FILE_NAME}}_re_id`') + row_count_json=$(bq query --use_legacy_sql=false --format json 'SELECT COUNT(ID) FROM `${{env.PROJECT_ID}}.${{env.REID_WITHOUT_QUERY_DATASET_ID}}.${{env.INPUT_FILE_NAME}}_re_id`') row_count=$(echo "$row_count_json" | jq -r '.[].f0_') if [[ "$row_count" == "$rc_orig" ]]; then echo "PASSED"; @@ -725,11 +725,11 @@ jobs: env: DATASET_ID: ${{ needs.generate-uuid.outputs.output5 }} PARQUET_DATASET_ID: ${{ needs.generate-uuid.outputs.output8 }} - REID_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} + REID_WITHOUT_QUERY_DATASET_ID: ${{ needs.generate-uuid.outputs.output10 }} run: | bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.DATASET_ID}} bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.PARQUET_DATASET_ID}} - bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.REID_DATASET_ID}} + bq rm -r -f -d ${{env.PROJECT_ID}}:${{env.REID_WITHOUT_QUERY_DATASET_ID}} - name: Clean up pub_sub file run: | diff --git a/README.md b/README.md index 5b1c9974..be9c858f 100644 --- a/README.md +++ b/README.md @@ -405,6 +405,7 @@ to validate de-identified results: #### Re-identification from BigQuery + 1. Export a SQL query to read and re-identify data from BigQuery. The sample provided below selects 10 records that match the query. ``` @@ -447,6 +448,7 @@ to validate de-identified results: stored in `reid_query.sql`. The re-identified results can be found in the BigQuery dataset (`dataset` parameter) with the name of the input table as the suffix. + The parameter `queryPath` is optional. If not passed, the pipeline will perform re-identification on the entire BigQuery table. ### Pipeline Parameters @@ -471,7 +473,7 @@ to validate de-identified results: | `recordDelimiter` | (Optional) Record delimiter. | INSPECT/DEID | | `columnDelimiter` | Column delimiter. Only required in case of a custom delimiter. | INSPECT/DEID | | `tableRef` | BigQuery table to export from in the form `:.`. | REID | -| `queryPath` | Query file for re-identification. | REID | +| `queryPath` | (Optional) Query file for re-identification. | REID | | `headers` | DLP table headers. Required for the JSONL file type. | INSPECT/DEID | | `numShardsPerDLPRequestBatching` | (Optional) Number of shards for DLP request batches. Can be used to control the parallelism of DLP requests. The default value is 100. | All | | `numberOfWorkerHarnessThreads` | (Optional) The number of threads per each worker harness process. | All |