From 02e515d2f23e7acd6648a9c9808956261df9bc00 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Wed, 19 Nov 2025 15:12:43 +0000 Subject: [PATCH 01/24] add spark runner documentation --- .github/workflows/bazel.yml | 8 +- .github/workflows/maven.yml | 7 +- examples/pipelinedp4j/README.md | 5 ++ examples/pipelinedp4j/WORKSPACE.bazel | 1 + examples/pipelinedp4j/beam/pom.xml | 79 +++++++++++++++++-- .../pipelinedp4j/examples/BUILD.bazel | 3 + 6 files changed, 94 insertions(+), 9 deletions(-) diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index a0b3a834..a5fcb16d 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -30,6 +30,9 @@ jobs: - name: Test C++ Examples working-directory: examples/cc run: bazelisk test --test_output=errors --test_verbose_timeout_warnings ... + - name: Run Beam example with SparkRunner + working-directory: examples/pipelinedp4j + run: bazelisk run beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:BeamExample -- --runner=SparkRunner --sparkMaster="[local]" --inputFilePath="$(pwd)/input.csv" --outputFilePath="$(pwd)/spark.txt" java-tests: name: Java Bazel Tests @@ -93,6 +96,9 @@ jobs: - name: Run Spark DataFrame Example working-directory: examples/pipelinedp4j run: bazelisk run spark/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:SparkDataFrameExample -- --inputFilePath="$(pwd)/input.csv" --outputFolder="$(pwd)/output" + - name: Run Spark Runner Example + working-directory: examples/pipelinedp4j + run: bazelisk run spark/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:SparkDataFrameExample -- --inputFilePath="$(pwd)/input.csv" --outputFolder="$(pwd)/output" zetasql-build: name: ZetaSQL Examples Build Test @@ -120,7 +126,7 @@ jobs: - name: Set up Python 3.10 uses: actions/setup-python@v4 with: - python-version: '3.10' + python-version: "3.10" - name: Install Python dependencies run: | python -m pip install --upgrade pip diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index 3b288b73..faf590f5 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -19,8 +19,8 @@ jobs: - name: Set up JDK 11 uses: actions/setup-java@v3 with: - java-version: '11' - distribution: 'temurin' + java-version: "11" + distribution: "temurin" - name: Install Maven run: sudo apt-get install -y maven - name: Create dummy input file for examples @@ -46,3 +46,6 @@ jobs: - name: Run Spark DataFrame Example working-directory: examples/pipelinedp4j/spark run: mvn compile exec:java -Dexec.mainClass=com.google.privacy.differentialprivacy.pipelinedp4j.examples.SparkDataFrameExample -Dexec.args="--inputFilePath=$(pwd)/../input.csv --outputFolder=output" + - name: Run Spark Runner Example + working-directory: examples/pipelinedp4j/spark + run: mvn compile exec:java -Dexec.mainClass=com.google.privacy.differentialprivacy.pipelinedp4j.examples.SparkDataFrameExample -Dexec.args="--inputFilePath=$(pwd)/../input.csv --outputFolder=output" diff --git a/examples/pipelinedp4j/README.md b/examples/pipelinedp4j/README.md index 023b407b..e042ecb6 100644 --- a/examples/pipelinedp4j/README.md +++ b/examples/pipelinedp4j/README.md @@ -159,6 +159,11 @@ For Spark the output is written to a folder and the result is stored in a file whose name starts with `part-00000`: `cat output/part-00000<...>` + +### Running with [Sparkrunner](https://beam.apache.org/documentation/runners/spark/) ( Beam ) + + + ### Running on Google Cloud Platform This section explains the examples on Google Cloud Platform (GCP). diff --git a/examples/pipelinedp4j/WORKSPACE.bazel b/examples/pipelinedp4j/WORKSPACE.bazel index 17e8271b..790b7841 100644 --- a/examples/pipelinedp4j/WORKSPACE.bazel +++ b/examples/pipelinedp4j/WORKSPACE.bazel @@ -86,6 +86,7 @@ maven_install( "com.google.protobuf:protobuf-kotlin:4.30.1", "com.google.errorprone:error_prone_annotations:2.38.0", "org.apache.beam:beam-runners-direct-java:%s" % BEAM_TAG, + "org.apache.beam:beam-runners-spark-3:%s" % BEAM_TAG, "org.apache.beam:beam-sdks-java-core:%s" % BEAM_TAG, "org.apache.beam:beam-sdks-java-extensions-avro:%s" % BEAM_TAG, "org.apache.beam:beam-sdks-java-extensions-protobuf:%s" % BEAM_TAG, diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 28ad351f..a839c0ec 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -17,8 +17,8 @@ --> + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 @@ -32,6 +32,7 @@ 2.63.0 + 3.5.0 false @@ -53,7 +54,7 @@ beam-sdks-java-core ${beam.version} - + @@ -73,16 +74,82 @@ - dataflow-runner + spark-runner org.apache.beam - beam-runners-google-cloud-dataflow-java + beam-runners-spark-3 ${beam.version} runtime + + + + + org.apache.maven.plugins + maven-shade-plugin + + false + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + com.google.guava:guava + + + + + + + package + + shade + + + true + shaded + + + com.google.common + relocated.com.google.common + + + + + + + + + + + + + + + + + spark-runner-embedeed + + + + org.apache.spark + spark-core_2.12 + ${spark.version} + + + + org.apache.spark + spark-streaming_2.12 + ${spark.version} + + - + \ No newline at end of file diff --git a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel index f5ca9c6a..09af8ce0 100644 --- a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel +++ b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel @@ -20,6 +20,9 @@ java_binary( main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", runtime_deps = [ "@maven//:org_apache_beam_beam_runners_direct_java", + "@maven//:org_apache_beam_beam_runners_spark_3", + + ], deps = [ "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", From f32709d8a5d6f1a80a4205dc0aa80155177ce5a2 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Wed, 19 Nov 2025 22:32:43 +0000 Subject: [PATCH 02/24] add maven.yaml support --- .github/workflows/maven.yml | 7 +- examples/pipelinedp4j/beam/pom.xml | 284 ++++++++++++++--------------- 2 files changed, 146 insertions(+), 145 deletions(-) diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index faf590f5..f7d51e8f 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -46,6 +46,9 @@ jobs: - name: Run Spark DataFrame Example working-directory: examples/pipelinedp4j/spark run: mvn compile exec:java -Dexec.mainClass=com.google.privacy.differentialprivacy.pipelinedp4j.examples.SparkDataFrameExample -Dexec.args="--inputFilePath=$(pwd)/../input.csv --outputFolder=output" + - name: Build Beam with Spark Runner + working-directory: examples/pipelinedp4j/beam + run: mvn package -Pspark-runner,spark-runner-embedeed - name: Run Spark Runner Example - working-directory: examples/pipelinedp4j/spark - run: mvn compile exec:java -Dexec.mainClass=com.google.privacy.differentialprivacy.pipelinedp4j.examples.SparkDataFrameExample -Dexec.args="--inputFilePath=$(pwd)/../input.csv --outputFolder=output" + working-directory: examples/pipelinedp4j/beam + run: java --add-opens=java.base/sun.nio.ch=ALL-UNNAMED -jar target/beam-1.0-SNAPSHOT-shaded.jar --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath=$(pwd)/../input.csv --outputFilePath=output-spark.txt diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index a839c0ec..ff6e4f67 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -1,155 +1,153 @@ - 4.0.0 - - +xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" +xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> +4.0.0 + + + com.google.privacy.differentialprivacy + pipelinedp4j + 1.0-SNAPSHOT + + +beam +jar + + + 2.63.0 + 3.5.0 + + + false + + UTF-8 + + + + + com.google.privacy.differentialprivacy - pipelinedp4j + common 1.0-SNAPSHOT - - - beam - jar - - - 2.63.0 - 3.5.0 - - - false - - UTF-8 - - - - - - com.google.privacy.differentialprivacy - common - 1.0-SNAPSHOT - - - - - org.apache.beam - beam-sdks-java-core - ${beam.version} - - - - - - direct-runner - - true - - - - - org.apache.beam - beam-runners-direct-java - ${beam.version} - runtime - - - - - - spark-runner - - - - org.apache.beam - beam-runners-spark-3 - ${beam.version} - runtime - - - - - - - org.apache.maven.plugins - maven-shade-plugin - - false - - - *:* - - META-INF/*.SF - META-INF/*.DSA - META-INF/*.RSA - com.google.guava:guava - - - - - - - package - - shade - - - true - shaded - - - com.google.common - relocated.com.google.common - - - - - - - - - - - - - - - - - spark-runner-embedeed - - - - org.apache.spark - spark-core_2.12 - ${spark.version} - - - - org.apache.spark - spark-streaming_2.12 - ${spark.version} - - - - + + + + + org.apache.beam + beam-sdks-java-core + ${beam.version} + + + + + + direct-runner + + true + + + + + org.apache.beam + beam-runners-direct-java + ${beam.version} + runtime + + + + + + spark-runner + + + + org.apache.beam + beam-runners-spark-3 + ${beam.version} + runtime + + + + + + + org.apache.maven.plugins + maven-shade-plugin + + false + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + + + + + + + package + + shade + + + true + shaded + + + + + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + + + + + + + + + + + + + + spark-runner-embedeed + + + + org.apache.spark + spark-core_2.12 + ${spark.version} + + + + org.apache.spark + spark-streaming_2.12 + ${spark.version} + + + + + \ No newline at end of file From daf611ef1d32d42d855e59872deca8d5fe841e96 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:13:34 +0000 Subject: [PATCH 03/24] add bazel support --- .github/workflows/bazel.yml | 2 +- examples/pipelinedp4j/MODULE.bazel | 6 +++++ examples/pipelinedp4j/WORKSPACE.bazel | 25 ++++++++++++++++++- .../pipelinedp4j/examples/BUILD.bazel | 4 ++- pipelinedp4j/MODULE.bazel | 6 +++++ 5 files changed, 40 insertions(+), 3 deletions(-) create mode 100644 examples/pipelinedp4j/MODULE.bazel create mode 100644 pipelinedp4j/MODULE.bazel diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index a5fcb16d..3fced511 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -98,7 +98,7 @@ jobs: run: bazelisk run spark/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:SparkDataFrameExample -- --inputFilePath="$(pwd)/input.csv" --outputFolder="$(pwd)/output" - name: Run Spark Runner Example working-directory: examples/pipelinedp4j - run: bazelisk run spark/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:SparkDataFrameExample -- --inputFilePath="$(pwd)/input.csv" --outputFolder="$(pwd)/output" + run: bazelisk run beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:BeamExample -- --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath="$(pwd)/input.csv" --outputFilePath="$(pwd)/output-sparkrunner.txt" zetasql-build: name: ZetaSQL Examples Build Test diff --git a/examples/pipelinedp4j/MODULE.bazel b/examples/pipelinedp4j/MODULE.bazel new file mode 100644 index 00000000..00bb1836 --- /dev/null +++ b/examples/pipelinedp4j/MODULE.bazel @@ -0,0 +1,6 @@ +############################################################################### +# Bazel now uses Bzlmod by default to manage external dependencies. +# Please consider migrating your external dependencies from WORKSPACE to MODULE.bazel. +# +# For more details, please check https://github.com/bazelbuild/bazel/issues/18958 +############################################################################### diff --git a/examples/pipelinedp4j/WORKSPACE.bazel b/examples/pipelinedp4j/WORKSPACE.bazel index 790b7841..22f080cd 100644 --- a/examples/pipelinedp4j/WORKSPACE.bazel +++ b/examples/pipelinedp4j/WORKSPACE.bazel @@ -72,6 +72,8 @@ BEAM_TAG = "2.63.0" SCALA_TAG = "2.13" +SCALA_SPARK_RUNNER_TAG = "2.12" + SCALA_LIBRARY_TAG = "%s.16" % SCALA_TAG SPARK_TAG = "3.5.5" @@ -100,6 +102,8 @@ maven_install( "com.fasterxml.jackson.module:jackson-module-scala_%s:%s" % (SCALA_TAG, JACKSON_TAG), "org.scala-lang:scala-library:%s" % SCALA_LIBRARY_TAG, "info.picocli:picocli:4.7.6", + # For Apache Spark Runner testing locally + "org.apache.spark:spark-streaming_%s:%s" % (2.12, SPARK_TAG), ], repositories = [ "https://jcenter.bintray.com", @@ -107,8 +111,27 @@ maven_install( "https://repo.maven.apache.org/maven2", "https://repo1.maven.org/maven2", ], + +) +maven_install( + name = "maven_spark_2_12", + artifacts = [ + # Beam SDKs (Needed here to resolve transitive deps correctly within this namespace) + + # Spark (Scala 2.12) - Strictly 2.12 + "org.apache.spark:spark-core_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), + "org.apache.spark:spark-streaming_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), + "org.apache.spark:spark-sql_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), + "org.scala-lang:scala-library:%s.18" % SCALA_SPARK_RUNNER_TAG, + "com.fasterxml.jackson.module:jackson-module-scala_%s:%s" % (SCALA_SPARK_RUNNER_TAG,JACKSON_TAG), + ], + repositories = [ + "https://jcenter.bintray.com", + "https://maven.google.com", + "https://repo.maven.apache.org/maven2", + "https://repo1.maven.org/maven2", + ], ) - # gRPC grpc_kt_repositories() diff --git a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel index 09af8ce0..de57ef20 100644 --- a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel +++ b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel @@ -21,7 +21,9 @@ java_binary( runtime_deps = [ "@maven//:org_apache_beam_beam_runners_direct_java", "@maven//:org_apache_beam_beam_runners_spark_3", - + # Enforce Scala 2.12 artifacts as the Spark Runner is compiled against Scala 2.12 + "@maven_spark_2_12//:org_apache_spark_spark_core_2_12", + "@maven_spark_2_12//:org_apache_spark_spark_streaming_2_12", ], deps = [ diff --git a/pipelinedp4j/MODULE.bazel b/pipelinedp4j/MODULE.bazel new file mode 100644 index 00000000..00bb1836 --- /dev/null +++ b/pipelinedp4j/MODULE.bazel @@ -0,0 +1,6 @@ +############################################################################### +# Bazel now uses Bzlmod by default to manage external dependencies. +# Please consider migrating your external dependencies from WORKSPACE to MODULE.bazel. +# +# For more details, please check https://github.com/bazelbuild/bazel/issues/18958 +############################################################################### From 95f27c8532d2ab308812210ec41fb8d501bb8c11 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:43:22 +0000 Subject: [PATCH 04/24] update documentation --- examples/pipelinedp4j/README.md | 42 +++++++++++++++++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/examples/pipelinedp4j/README.md b/examples/pipelinedp4j/README.md index e042ecb6..59ec1e80 100644 --- a/examples/pipelinedp4j/README.md +++ b/examples/pipelinedp4j/README.md @@ -162,7 +162,49 @@ output/part-00000<...>` ### Running with [Sparkrunner](https://beam.apache.org/documentation/runners/spark/) ( Beam ) +You can execute PipelineDP4j Beam examples using the Spark Runner. This allows your Beam pipeline to run on a Spark cluster or a locally simulated Spark environment. +# Build and submit to a Spark Cluster + +To run the Beam example with SparkRunner, you need to use the `spark-runner` profile and add the `--runner=SparkRunner` parameter. + + +First, build the package with the `spark-runner` profile: + +```shell +mvn clean package -Pspark-runner +``` + +Then submit with Spark CLI: + +```shell +spark-submit \ + --class com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample \ + --master spark://localhost:7077 \ + target/beam-1.0-SNAPSHOT-shaded.jar \ + --runner=SparkRunner \ + --inputFilePath=/netflix_data.csv \ + --outputFilePath=output-spark-runner.txt +``` + +View the results with `cat output-spark-runner.txt`. + +# Build and execute locallcy + +You can also run PipelineDP4j with the Spark Runner locally. +When running locally, the dependencies usually provided by the Spark Cluster are missing. Therefore, you must use the `spark-runner-embedeed` profile in addition to `spark-runner` to bundle the necessary Spark dependencies into your JAR. + +```shell +mvn clean package -Pspark-runner,spark-runner-embedeed +``` + +Then execute the JAR file + +```shell +java --add-opens=java.base/sun.nio.ch=ALL-UNNAMED -jar target/beam-1.0-SNAPSHOT-shaded.jar --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath=/netflix_data.csv --outputFilePath=output-spark-runner.txt +``` + +View the results with `cat output-spark-runner.txt`. ### Running on Google Cloud Platform From 8c96aaf1e6be27bdd4d8ab4e86994ab498aa57f8 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:45:51 +0000 Subject: [PATCH 05/24] remove unnecessary job --- .github/workflows/bazel.yml | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index 3fced511..772a6b4d 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -30,10 +30,6 @@ jobs: - name: Test C++ Examples working-directory: examples/cc run: bazelisk test --test_output=errors --test_verbose_timeout_warnings ... - - name: Run Beam example with SparkRunner - working-directory: examples/pipelinedp4j - run: bazelisk run beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:BeamExample -- --runner=SparkRunner --sparkMaster="[local]" --inputFilePath="$(pwd)/input.csv" --outputFilePath="$(pwd)/spark.txt" - java-tests: name: Java Bazel Tests runs-on: ubuntu-latest @@ -96,7 +92,7 @@ jobs: - name: Run Spark DataFrame Example working-directory: examples/pipelinedp4j run: bazelisk run spark/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:SparkDataFrameExample -- --inputFilePath="$(pwd)/input.csv" --outputFolder="$(pwd)/output" - - name: Run Spark Runner Example + - name: Run Beam Spark Runner Example working-directory: examples/pipelinedp4j run: bazelisk run beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:BeamExample -- --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath="$(pwd)/input.csv" --outputFilePath="$(pwd)/output-sparkrunner.txt" From d862d44ec05c418377a6bb4152e68931de069bac Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:48:34 +0000 Subject: [PATCH 06/24] Revert removal of empty line --- .github/workflows/bazel.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index 772a6b4d..90d766e0 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -30,6 +30,7 @@ jobs: - name: Test C++ Examples working-directory: examples/cc run: bazelisk test --test_output=errors --test_verbose_timeout_warnings ... + java-tests: name: Java Bazel Tests runs-on: ubuntu-latest From 883c0bf03ec7893dfa8bcb0aeddfdc9578a35988 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:50:12 +0000 Subject: [PATCH 07/24] delete unnecessary bazel file --- examples/pipelinedp4j/MODULE.bazel | 6 ------ 1 file changed, 6 deletions(-) delete mode 100644 examples/pipelinedp4j/MODULE.bazel diff --git a/examples/pipelinedp4j/MODULE.bazel b/examples/pipelinedp4j/MODULE.bazel deleted file mode 100644 index 00bb1836..00000000 --- a/examples/pipelinedp4j/MODULE.bazel +++ /dev/null @@ -1,6 +0,0 @@ -############################################################################### -# Bazel now uses Bzlmod by default to manage external dependencies. -# Please consider migrating your external dependencies from WORKSPACE to MODULE.bazel. -# -# For more details, please check https://github.com/bazelbuild/bazel/issues/18958 -############################################################################### From d288a366a01748c71a05344cf17aac26e272d492 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 12:51:59 +0000 Subject: [PATCH 08/24] remove unnecessary bazel file --- pipelinedp4j/MODULE.bazel | 6 ------ 1 file changed, 6 deletions(-) delete mode 100644 pipelinedp4j/MODULE.bazel diff --git a/pipelinedp4j/MODULE.bazel b/pipelinedp4j/MODULE.bazel deleted file mode 100644 index 00bb1836..00000000 --- a/pipelinedp4j/MODULE.bazel +++ /dev/null @@ -1,6 +0,0 @@ -############################################################################### -# Bazel now uses Bzlmod by default to manage external dependencies. -# Please consider migrating your external dependencies from WORKSPACE to MODULE.bazel. -# -# For more details, please check https://github.com/bazelbuild/bazel/issues/18958 -############################################################################### From 5881c88a8b4c7f9bde4f4fe38a23b9cb9a8dcb3c Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:03:55 +0000 Subject: [PATCH 09/24] fix indent --- examples/pipelinedp4j/README.md | 2 +- examples/pipelinedp4j/WORKSPACE.bazel | 1 + examples/pipelinedp4j/beam/pom.xml | 275 ++++++++++++++------------ 3 files changed, 146 insertions(+), 132 deletions(-) diff --git a/examples/pipelinedp4j/README.md b/examples/pipelinedp4j/README.md index 59ec1e80..2d9f5b2d 100644 --- a/examples/pipelinedp4j/README.md +++ b/examples/pipelinedp4j/README.md @@ -180,7 +180,7 @@ Then submit with Spark CLI: ```shell spark-submit \ --class com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample \ - --master spark://localhost:7077 \ + --master spark://: \ target/beam-1.0-SNAPSHOT-shaded.jar \ --runner=SparkRunner \ --inputFilePath=/netflix_data.csv \ diff --git a/examples/pipelinedp4j/WORKSPACE.bazel b/examples/pipelinedp4j/WORKSPACE.bazel index 22f080cd..3b94c0f5 100644 --- a/examples/pipelinedp4j/WORKSPACE.bazel +++ b/examples/pipelinedp4j/WORKSPACE.bazel @@ -132,6 +132,7 @@ maven_install( "https://repo1.maven.org/maven2", ], ) + # gRPC grpc_kt_repositories() diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index ff6e4f67..4032d0ce 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -17,137 +17,150 @@ limitations under the License. --> -4.0.0 - - - com.google.privacy.differentialprivacy - pipelinedp4j - 1.0-SNAPSHOT - - -beam -jar - - - 2.63.0 - 3.5.0 - - - false - - UTF-8 - - - - - + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> + 4.0.0 + + com.google.privacy.differentialprivacy - common + pipelinedp4j 1.0-SNAPSHOT - - - - - org.apache.beam - beam-sdks-java-core - ${beam.version} - - - - - - direct-runner - - true - - - - - org.apache.beam - beam-runners-direct-java - ${beam.version} - runtime - - - - - - spark-runner - - - - org.apache.beam - beam-runners-spark-3 - ${beam.version} - runtime - - - - - - - org.apache.maven.plugins - maven-shade-plugin - - false - - - *:* - - META-INF/*.SF - META-INF/*.DSA - META-INF/*.RSA - - - - - - - package - - shade - - - true - shaded - - - - - com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample - - - - - - - - - - - - - - spark-runner-embedeed - - - - org.apache.spark - spark-core_2.12 - ${spark.version} - - - - org.apache.spark - spark-streaming_2.12 - ${spark.version} - - - - - + + + beam + jar + + + 2.63.0 + 3.5.0 + + + false + + UTF-8 + + + + + + com.google.privacy.differentialprivacy + common + 1.0-SNAPSHOT + + + + + org.apache.beam + beam-sdks-java-core + ${beam.version} + + + + + + direct-runner + + true + + + + + org.apache.beam + beam-runners-direct-java + ${beam.version} + runtime + + + + + + + + dataflow-runner + + + + org.apache.beam + beam-runners-google-cloud-dataflow-java + ${beam.version} + runtime + + + + spark-runner + + + + org.apache.beam + beam-runners-spark-3 + ${beam.version} + runtime + + + + + + + org.apache.maven.plugins + maven-shade-plugin + + false + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + + + + + + + package + + shade + + + true + shaded + + + + + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + + + + + + + + + + + + + + spark-runner-embedeed + + + + org.apache.spark + spark-core_2.12 + ${spark.version} + + + + org.apache.spark + spark-streaming_2.12 + ${spark.version} + + + + + \ No newline at end of file From 78d56ff70ddb95b2ba6332202baec871e547de84 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:06:38 +0000 Subject: [PATCH 10/24] fix profile typo --- examples/pipelinedp4j/beam/pom.xml | 118 ++++++++++++++--------------- 1 file changed, 58 insertions(+), 60 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 4032d0ce..82e437b5 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -73,75 +73,73 @@ limitations under the License. - - - dataflow-runner - - - - org.apache.beam - beam-runners-google-cloud-dataflow-java - ${beam.version} - runtime - - - - spark-runner + dataflow-runner org.apache.beam - beam-runners-spark-3 + beam-runners-google-cloud-dataflow-java ${beam.version} runtime - - - - - org.apache.maven.plugins - maven-shade-plugin - - false - - - *:* - - META-INF/*.SF - META-INF/*.DSA - META-INF/*.RSA - - - - - - - package - - shade - - - true - shaded - - - - - com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample - - - - - - - - - - + + spark-runner + + + + org.apache.beam + beam-runners-spark-3 + ${beam.version} + runtime + + + + + + + org.apache.maven.plugins + maven-shade-plugin + + false + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + + + + + + + package + + shade + + + true + shaded + + + + + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + + + + + + + + + + From 33cab369a2631ba09896e98751358252e85d57ba Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:09:52 +0000 Subject: [PATCH 11/24] fix type ( double quote instead of simple quote) --- .github/workflows/bazel.yml | 2 +- .github/workflows/maven.yml | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/bazel.yml b/.github/workflows/bazel.yml index 90d766e0..eb4bd7bb 100644 --- a/.github/workflows/bazel.yml +++ b/.github/workflows/bazel.yml @@ -123,7 +123,7 @@ jobs: - name: Set up Python 3.10 uses: actions/setup-python@v4 with: - python-version: "3.10" + python-version: '3.10' - name: Install Python dependencies run: | python -m pip install --upgrade pip diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index f7d51e8f..26fbfd29 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -19,8 +19,8 @@ jobs: - name: Set up JDK 11 uses: actions/setup-java@v3 with: - java-version: "11" - distribution: "temurin" + java-version: '11' + distribution: 'temurin' - name: Install Maven run: sudo apt-get install -y maven - name: Create dummy input file for examples From b3622c273371aca1dc3b5b7b5be271d98f65213c Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:12:07 +0000 Subject: [PATCH 12/24] indent xml file --- examples/pipelinedp4j/beam/pom.xml | 24 ++++++++++++------------ 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 82e437b5..d05cdb4e 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -1,24 +1,24 @@ + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://www.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 From 07eb1f22996394a00d15138fc03d4c4a1b2ba1b4 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:13:02 +0000 Subject: [PATCH 13/24] revert schema location --- examples/pipelinedp4j/beam/pom.xml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index d05cdb4e..77c49773 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -17,8 +17,8 @@ --> + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 From a17030bd0cfeff53b61182d81e940f05b375ce32 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:15:13 +0000 Subject: [PATCH 14/24] fix missing tag --- examples/pipelinedp4j/beam/pom.xml | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 77c49773..1720d017 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -85,8 +85,9 @@ + spark-runner - + org.apache.beam @@ -161,4 +162,4 @@ - \ No newline at end of file + From f9649294868a50ae08644d37d899676bf313a2e5 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 20 Nov 2025 13:17:39 +0000 Subject: [PATCH 15/24] reformat pom.xml --- examples/pipelinedp4j/beam/pom.xml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 1720d017..110b4f73 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -17,8 +17,8 @@ --> + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 @@ -54,7 +54,7 @@ beam-sdks-java-core ${beam.version} - + From 86e25de8eae745df2f6d10ddf1d06a43550090a4 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 29 Jan 2026 14:46:32 +0000 Subject: [PATCH 16/24] nit: typo, renaming, indentation --- .github/workflows/maven.yml | 4 +- examples/pipelinedp4j/README.md | 6 +- examples/pipelinedp4j/beam/pom.xml | 212 +++++++++--------- .../pipelinedp4j/examples/BUILD.bazel | 36 +++ 4 files changed, 147 insertions(+), 111 deletions(-) create mode 100644 examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index 26fbfd29..284a8ca0 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -49,6 +49,6 @@ jobs: - name: Build Beam with Spark Runner working-directory: examples/pipelinedp4j/beam run: mvn package -Pspark-runner,spark-runner-embedeed - - name: Run Spark Runner Example + - name: Run Beam Example with Spark Runner working-directory: examples/pipelinedp4j/beam - run: java --add-opens=java.base/sun.nio.ch=ALL-UNNAMED -jar target/beam-1.0-SNAPSHOT-shaded.jar --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath=$(pwd)/../input.csv --outputFilePath=output-spark.txt + run: java --add-opens=java.base/sun.nio.ch=ALL-UNNAMED -jar target/beam-1.0-SNAPSHOT-shaded.jar --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath=../input.csv --outputFilePath=output-spark.txt diff --git a/examples/pipelinedp4j/README.md b/examples/pipelinedp4j/README.md index 2d9f5b2d..aef901e9 100644 --- a/examples/pipelinedp4j/README.md +++ b/examples/pipelinedp4j/README.md @@ -160,11 +160,11 @@ result is stored in a file whose name starts with `part-00000`: `cat output/part-00000<...>` -### Running with [Sparkrunner](https://beam.apache.org/documentation/runners/spark/) ( Beam ) +#### Running with [Sparkrunner](https://beam.apache.org/documentation/runners/spark/) (Beam ) You can execute PipelineDP4j Beam examples using the Spark Runner. This allows your Beam pipeline to run on a Spark cluster or a locally simulated Spark environment. -# Build and submit to a Spark Cluster +##### Build and submit to a Spark Cluster To run the Beam example with SparkRunner, you need to use the `spark-runner` profile and add the `--runner=SparkRunner` parameter. @@ -189,7 +189,7 @@ spark-submit \ View the results with `cat output-spark-runner.txt`. -# Build and execute locallcy +# Build and execute locally You can also run PipelineDP4j with the Spark Runner locally. When running locally, the dependencies usually provided by the Spark Cluster are missing. Therefore, you must use the `spark-runner-embedeed` profile in addition to `spark-runner` to bundle the necessary Spark dependencies into your JAR. diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 110b4f73..dec2f71a 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -1,93 +1,93 @@ - 4.0.0 - - +xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" +xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> +4.0.0 + + + com.google.privacy.differentialprivacy + pipelinedp4j + 1.0-SNAPSHOT + + +beam +jar + + + 2.63.0 + 3.5.0 + + + false + + UTF-8 + + + + + com.google.privacy.differentialprivacy - pipelinedp4j + common 1.0-SNAPSHOT - - - beam - jar - - - 2.63.0 - 3.5.0 - - - false - - UTF-8 - - - - - - com.google.privacy.differentialprivacy - common - 1.0-SNAPSHOT - - - - - org.apache.beam - beam-sdks-java-core - ${beam.version} - - - - - - direct-runner - - true - - - - - org.apache.beam - beam-runners-direct-java - ${beam.version} - runtime - - - - - - dataflow-runner - - - - org.apache.beam - beam-runners-google-cloud-dataflow-java - ${beam.version} - runtime - - - - + + + + + org.apache.beam + beam-sdks-java-core + ${beam.version} + + + + + + direct-runner + + true + + + + + org.apache.beam + beam-runners-direct-java + ${beam.version} + runtime + + + + + + dataflow-runner + + + + org.apache.beam + beam-runners-google-cloud-dataflow-java + ${beam.version} + runtime + + + + spark-runner - + org.apache.beam @@ -96,7 +96,7 @@ runtime - + @@ -124,42 +124,42 @@ true shaded - + - - com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> + + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer" /> - + - - - - spark-runner-embedeed - - - - org.apache.spark - spark-core_2.12 - ${spark.version} - - - - org.apache.spark - spark-streaming_2.12 - ${spark.version} - - - - - + + + + spark-runner-embedeed + + + + org.apache.spark + spark-core_2.12 + ${spark.version} + + + + org.apache.spark + spark-streaming_2.12 + ${spark.version} + + + + + diff --git a/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel b/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel new file mode 100644 index 00000000..de57ef20 --- /dev/null +++ b/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel @@ -0,0 +1,36 @@ +# Copyright 2024 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +java_binary( + name = "BeamExample", + srcs = [ + "BeamExample.java", + ], + main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", + runtime_deps = [ + "@maven//:org_apache_beam_beam_runners_direct_java", + "@maven//:org_apache_beam_beam_runners_spark_3", + # Enforce Scala 2.12 artifacts as the Spark Runner is compiled against Scala 2.12 + "@maven_spark_2_12//:org_apache_spark_spark_core_2_12", + "@maven_spark_2_12//:org_apache_spark_spark_streaming_2_12", + + ], + deps = [ + "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", + "@com_google_privacy_differentialprivacy_pipelinedp4j//main/com/google/privacy/differentialprivacy/pipelinedp4j/api", + "@maven//:org_apache_beam_beam_sdks_java_core", + "@maven//:org_apache_beam_beam_sdks_java_extensions_avro", + "@maven//:org_jetbrains_kotlin_kotlin_stdlib", + ], +) From 9b5d30c0dbfa9e710052f52fa735038992ab61be Mon Sep 17 00:00:00 2001 From: leroyjb Date: Thu, 29 Jan 2026 14:48:43 +0000 Subject: [PATCH 17/24] nit: typo, renaming, indentation --- examples/pipelinedp4j/WORKSPACE.bazel | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/examples/pipelinedp4j/WORKSPACE.bazel b/examples/pipelinedp4j/WORKSPACE.bazel index 3b94c0f5..9967bc18 100644 --- a/examples/pipelinedp4j/WORKSPACE.bazel +++ b/examples/pipelinedp4j/WORKSPACE.bazel @@ -102,7 +102,7 @@ maven_install( "com.fasterxml.jackson.module:jackson-module-scala_%s:%s" % (SCALA_TAG, JACKSON_TAG), "org.scala-lang:scala-library:%s" % SCALA_LIBRARY_TAG, "info.picocli:picocli:4.7.6", - # For Apache Spark Runner testing locally + # For Apache Spark Runner running locally "org.apache.spark:spark-streaming_%s:%s" % (2.12, SPARK_TAG), ], repositories = [ From 46f231f354a1ca13606dcf6973a4d2f8fef68e90 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Tue, 21 Apr 2026 20:34:11 +0000 Subject: [PATCH 18/24] refactor documentation for Spark Runner with Beam --- examples/pipelinedp4j/README.md | 109 +++++++++++++----- examples/pipelinedp4j/WORKSPACE.bazel | 41 +++++-- examples/pipelinedp4j/beam/pom.xml | 43 ++++++- .../pipelinedp4j/examples/BUILD.bazel | 49 +++++++- .../pipelinedp4j/examples/BUILD.bazel | 36 ------ 5 files changed, 193 insertions(+), 85 deletions(-) delete mode 100644 examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel diff --git a/examples/pipelinedp4j/README.md b/examples/pipelinedp4j/README.md index aef901e9..22e17fd7 100644 --- a/examples/pipelinedp4j/README.md +++ b/examples/pipelinedp4j/README.md @@ -13,6 +13,9 @@ Next, explore the following options for building and running the example: Finally, delve into the [code walkthrough](#code-walkthrough) for a comprehensive understanding of how PipelineDP4j was employed to solve the task. +> [!TIP] +> If you are looking for specific `mvn` and `bazel` commands, you should have a look at the CI/CD files in the `.github` folder. You will find useful, up-to-date examples of `mvn` and `bazel` commands used in our automated workflows. + ## Problem statement This example demonstrates how to compute differentially private statistics on a @@ -159,37 +162,7 @@ For Spark the output is written to a folder and the result is stored in a file whose name starts with `part-00000`: `cat output/part-00000<...>` - -#### Running with [Sparkrunner](https://beam.apache.org/documentation/runners/spark/) (Beam ) - -You can execute PipelineDP4j Beam examples using the Spark Runner. This allows your Beam pipeline to run on a Spark cluster or a locally simulated Spark environment. - -##### Build and submit to a Spark Cluster - -To run the Beam example with SparkRunner, you need to use the `spark-runner` profile and add the `--runner=SparkRunner` parameter. - - -First, build the package with the `spark-runner` profile: - -```shell -mvn clean package -Pspark-runner -``` - -Then submit with Spark CLI: - -```shell -spark-submit \ - --class com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample \ - --master spark://: \ - target/beam-1.0-SNAPSHOT-shaded.jar \ - --runner=SparkRunner \ - --inputFilePath=/netflix_data.csv \ - --outputFilePath=output-spark-runner.txt -``` - -View the results with `cat output-spark-runner.txt`. - -# Build and execute locally +### Running locally with SparkRunner You can also run PipelineDP4j with the Spark Runner locally. When running locally, the dependencies usually provided by the Spark Cluster are missing. Therefore, you must use the `spark-runner-embedeed` profile in addition to `spark-runner` to bundle the necessary Spark dependencies into your JAR. @@ -283,6 +256,80 @@ to ensure your environment is correctly set up. Then do the following. After finish you can inspect the result on GCP bucket. + +#### Running on Dataproc (Spark) with [SparkRunner](https://beam.apache.org/documentation/runners/spark/) (Beam) + +> [!WARNING] +> SparkRunner use a 2.12 scala version that can lead to conflict with your existing code of your Spark version. Be cautious when you build you jar to include the correct version of your libraries, compatible with the SparkRunner. + +You can execute PipelineDP4j Beam examples using the Spark Runner. This allows your Beam pipeline to run on a Spark cluster or a locally simulated Spark environment. + +##### Build with Maven + +To run the Beam example with SparkRunner, you need to use the `spark-runner` profile and add the `--runner=SparkRunner` parameter. + + +First, build the package with the `spark-runner` profile: + +```shell +mvn clean package -Pspark-runner +``` + +Then submit the Spark job: + +```shell +gcloud dataproc batches submit spark --version 1.2 --region=us-central1 --jars=beam/target/beam-1.0-SNAPSHOT-shaded.jar --class=com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample --deps-bucket=gs:// -- --runner=SparkRunner --inputFilePath=gs:///netflix_data.csv --outputFilePath=gs:///output +``` + +##### Build with Bazel and library sources + +To run the Beam example with SparkRunner, please ensure you first follow the build commands outlined in the [Running using Bazel and library sources](#running-using-bazel-and-library-sources) section. + +Once built, submit the Spark job: + +Then submit the Spark job: + +```shell +gcloud dataproc batches submit spark --version 1.2 --region=us-central1 --jars=bazel-bin/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BeamExampleSparkRunnerCluster_shaded.jar --class=com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample --deps-bucket=gs:// -- --runner=SparkRunner --inputFilePath=gs:///netflix_data.csv --outputFilePath=gs:///output +``` + +After finish you can inspect the result on GCP bucket. + +##### Build and execute locally + +To test your pipeline before submitting it to Spark, you can run it locally using the SparkRunner. Note that when running locally, the dependencies typically provided by the Spark cluster will be missing. + +###### Build and execute locally with Maven + + Therefore, you must use the `spark-runner-embedeed` profile in addition to `spark-runner` to bundle the necessary Spark dependencies into your JAR. + +```shell +mvn clean package -Pspark-runner,spark-runner-embeddeed +``` + +Then execute the JAR file + +```shell +java --add-opens=java.base/sun.nio.ch=ALL-UNNAMED -jar beam/target/beam-1.0-SNAPSHOT-shaded.jar --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath=/netflix_data.csv --outputFilePath=output-spark-runner.txt +``` + +View the results with `cat output-spark-runner.txt`. + +###### Build and execute locally with Bazel + +To run the Beam example with SparkRunner, please ensure you first follow the build commands outlined in the [Running using Bazel and library sources](#running-using-bazel-and-library-sources) section. + +Then execute the Jar file : + +```shell +bazelisk run beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:BeamExampleSparkRunnerLocal -- --runner=SparkRunner --sparkMaster="local[*]" --inputFilePath="$(pwd)/input.csv" --outputFilePath="$(pwd)/output-sparkrunner.txt" + +``` + +View the results with `cat output-spark-runner.txt`. + + + ## Code walkthrough Let's deep into details how code for computing DP statistics is organized. diff --git a/examples/pipelinedp4j/WORKSPACE.bazel b/examples/pipelinedp4j/WORKSPACE.bazel index 9967bc18..b1e8f132 100644 --- a/examples/pipelinedp4j/WORKSPACE.bazel +++ b/examples/pipelinedp4j/WORKSPACE.bazel @@ -68,7 +68,7 @@ load( ) # Maven -BEAM_TAG = "2.63.0" +BEAM_TAG = "2.72.0" SCALA_TAG = "2.13" @@ -76,7 +76,7 @@ SCALA_SPARK_RUNNER_TAG = "2.12" SCALA_LIBRARY_TAG = "%s.16" % SCALA_TAG -SPARK_TAG = "3.5.5" +SPARK_TAG = "3.5.8" JACKSON_TAG = "2.18.3" @@ -103,8 +103,9 @@ maven_install( "org.scala-lang:scala-library:%s" % SCALA_LIBRARY_TAG, "info.picocli:picocli:4.7.6", # For Apache Spark Runner running locally - "org.apache.spark:spark-streaming_%s:%s" % (2.12, SPARK_TAG), + "org.apache.spark:spark-streaming_%s:%s" % (SCALA_TAG, SPARK_TAG), ], + repositories = [ "https://jcenter.bintray.com", "https://maven.google.com", @@ -116,14 +117,17 @@ maven_install( maven_install( name = "maven_spark_2_12", artifacts = [ - # Beam SDKs (Needed here to resolve transitive deps correctly within this namespace) + "org.apache.beam:beam-runners-spark-3:%s" % BEAM_TAG, + "org.apache.beam:beam-sdks-java-core:%s" % BEAM_TAG, + "org.apache.beam:beam-sdks-java-extensions-avro:%s" % BEAM_TAG, - # Spark (Scala 2.12) - Strictly 2.12 - "org.apache.spark:spark-core_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), - "org.apache.spark:spark-streaming_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), - "org.apache.spark:spark-sql_%s:%s" % (SCALA_SPARK_RUNNER_TAG,SPARK_TAG), - "org.scala-lang:scala-library:%s.18" % SCALA_SPARK_RUNNER_TAG, - "com.fasterxml.jackson.module:jackson-module-scala_%s:%s" % (SCALA_SPARK_RUNNER_TAG,JACKSON_TAG), + # Spark et Jackson strictement en 2.12 + "org.apache.spark:spark-core_2.12:%s" % SPARK_TAG, + "org.apache.spark:spark-streaming_2.12:%s" % SPARK_TAG, + "org.apache.spark:spark-sql_2.12:%s" % SPARK_TAG, + "com.fasterxml.jackson.module:jackson-module-scala_2.12:%s" % JACKSON_TAG, + "org.scala-lang:scala-library:2.12.18", + "org.jetbrains.kotlin:kotlin-stdlib:1.9.0", ], repositories = [ "https://jcenter.bintray.com", @@ -148,3 +152,20 @@ local_repository( name = "com_google_privacy_differentialprivacy_pipelinedp4j", path = "../../pipelinedp4j", ) + +# Uses Jarjar to shade JARs for the Apache Beam Spark Runner +load("@bazel_tools//tools/build_defs/repo:http.bzl", "http_archive") + +http_archive( + name = "bazel_jar_jar", + sha256 = "ac11a42b7e3de37dae6ebd54ea5c1e306a17af12004b7f48fc7a46c451fd94b9", + strip_prefix = "bazel_jar_jar-0.1.15", + url = "https://github.com/bazeltools/bazel_jar_jar/releases/download/v0.1.15/bazel_jar_jar-v0.1.15.tar.gz", +) + +load( + "@bazel_jar_jar//:jar_jar.bzl", + "jar_jar_repositories", +) + +jar_jar_repositories() \ No newline at end of file diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index dec2f71a..2a8fd7bd 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -31,8 +31,8 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xs jar - 2.63.0 - 3.5.0 + 2.72.0 + 3.5.8 false @@ -95,6 +95,12 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xs ${beam.version} runtime + + + com.google.cloud.bigdataoss + gcsio + 3.1.1 + @@ -106,6 +112,9 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xs false + *:* META-INF/*.SF @@ -122,18 +131,39 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xs shade + true + shaded + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + + + + + com.google.common + custom.guava.cache.com.google.common + + + @@ -144,8 +174,13 @@ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xs - spark-runner-embedeed - + spark-runner-embeddeed + org.apache.spark diff --git a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel index de57ef20..d9089293 100644 --- a/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel +++ b/examples/pipelinedp4j/beam/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel @@ -12,6 +12,8 @@ # See the License for the specific language governing permissions and # limitations under the License. +load("@bazel_jar_jar//:jar_jar.bzl", "jar_jar") + java_binary( name = "BeamExample", srcs = [ @@ -19,12 +21,44 @@ java_binary( ], main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", runtime_deps = [ - "@maven//:org_apache_beam_beam_runners_direct_java", - "@maven//:org_apache_beam_beam_runners_spark_3", - # Enforce Scala 2.12 artifacts as the Spark Runner is compiled against Scala 2.12 + "@maven//:org_apache_beam_beam_runners_direct_java", + ], + deps = [ + "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", + "@com_google_privacy_differentialprivacy_pipelinedp4j//main/com/google/privacy/differentialprivacy/pipelinedp4j/api", + "@maven//:org_apache_beam_beam_sdks_java_core", + "@maven//:org_apache_beam_beam_sdks_java_extensions_avro", + "@maven//:org_jetbrains_kotlin_kotlin_stdlib", + ], +) +# --- Target SPARK RUNNER LOCAL (Scala 2.12) --- +java_binary( + name = "BeamExampleSparkRunnerLocal", + srcs = ["BeamExample.java"], + main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", + runtime_deps = [ + "@maven_spark_2_12//:org_apache_beam_beam_runners_spark_3", "@maven_spark_2_12//:org_apache_spark_spark_core_2_12", "@maven_spark_2_12//:org_apache_spark_spark_streaming_2_12", - + "@maven_spark_2_12//:com_fasterxml_jackson_module_jackson_module_scala_2_12", + ], + deps = [ + "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", + "@com_google_privacy_differentialprivacy_pipelinedp4j//main/com/google/privacy/differentialprivacy/pipelinedp4j/api", + "@maven_spark_2_12//:org_apache_beam_beam_sdks_java_core", + "@maven_spark_2_12//:org_apache_beam_beam_sdks_java_extensions_avro", + "@maven_spark_2_12//:org_jetbrains_kotlin_kotlin_stdlib", + ], +) + +# Target for cluster deployment (equivalent to the spark-runner profile from Maven) +# Do NOT include Spark in the runtime dependencies, as the cluster already provides it. +java_binary( + name = "BeamExampleSparkRunnerCluster", + srcs = ["BeamExample.java"], + main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", + runtime_deps = [ + "@maven//:org_apache_beam_beam_runners_spark_3", ], deps = [ "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", @@ -34,3 +68,10 @@ java_binary( "@maven//:org_jetbrains_kotlin_kotlin_stdlib", ], ) +# Shade the jar: Relocate the Guava package to avoid dependency conflicts. +# The Spark runner brings in its own version of Guava compiled for Scala 2.12. +jar_jar( + name = "BeamExampleSparkRunnerCluster_shaded", + input_jar = ":BeamExampleSparkRunnerCluster_deploy.jar", + inline_rules = ["rule com.google.common.** custom.guava.cache.com.google.common.@1"], +) \ No newline at end of file diff --git a/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel b/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel deleted file mode 100644 index de57ef20..00000000 --- a/examples/pipelinedp4j/beam/target/classes/com/google/privacy/differentialprivacy/pipelinedp4j/examples/BUILD.bazel +++ /dev/null @@ -1,36 +0,0 @@ -# Copyright 2024 Google LLC -# -# Licensed under the Apache License, Version 2.0 (the "License"); -# you may not use this file except in compliance with the License. -# You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. - -java_binary( - name = "BeamExample", - srcs = [ - "BeamExample.java", - ], - main_class = "com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample", - runtime_deps = [ - "@maven//:org_apache_beam_beam_runners_direct_java", - "@maven//:org_apache_beam_beam_runners_spark_3", - # Enforce Scala 2.12 artifacts as the Spark Runner is compiled against Scala 2.12 - "@maven_spark_2_12//:org_apache_spark_spark_core_2_12", - "@maven_spark_2_12//:org_apache_spark_spark_streaming_2_12", - - ], - deps = [ - "//common/src/main/java/com/google/privacy/differentialprivacy/pipelinedp4j/examples:common", - "@com_google_privacy_differentialprivacy_pipelinedp4j//main/com/google/privacy/differentialprivacy/pipelinedp4j/api", - "@maven//:org_apache_beam_beam_sdks_java_core", - "@maven//:org_apache_beam_beam_sdks_java_extensions_avro", - "@maven//:org_jetbrains_kotlin_kotlin_stdlib", - ], -) From 159b4f8e6a274e849d7a98b887103b221511df9a Mon Sep 17 00:00:00 2001 From: leroyjb Date: Tue, 21 Apr 2026 20:46:23 +0000 Subject: [PATCH 19/24] reformat pom.xml --- examples/pipelinedp4j/beam/pom.xml | 316 ++++++++++++++--------------- 1 file changed, 147 insertions(+), 169 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 2a8fd7bd..be7bb224 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -1,7 +1,5 @@ - - - - -4.0.0 - - - com.google.privacy.differentialprivacy - pipelinedp4j - 1.0-SNAPSHOT - - -beam -jar - - - 2.72.0 - 3.5.8 - - - false - - UTF-8 - - - - - +limitations under the License. --> + + 4.0.0 + com.google.privacy.differentialprivacy - common + pipelinedp4j 1.0-SNAPSHOT - - - - - org.apache.beam - beam-sdks-java-core - ${beam.version} - - - - - - direct-runner - - true - - - - - org.apache.beam - beam-runners-direct-java - ${beam.version} - runtime - - - - - - dataflow-runner - - - - org.apache.beam - beam-runners-google-cloud-dataflow-java - ${beam.version} - runtime - - - - - spark-runner - - - - org.apache.beam - beam-runners-spark-3 - ${beam.version} - runtime - + + beam + jar + + 2.72.0 + 3.5.8 + + false + + UTF-8 + + + + + com.google.privacy.differentialprivacy + common + 1.0-SNAPSHOT + + + + org.apache.beam + beam-sdks-java-core + ${beam.version} + + + + + direct-runner + + true + + + + + org.apache.beam + beam-runners-direct-java + ${beam.version} + runtime + + + + + dataflow-runner + + + + org.apache.beam + beam-runners-google-cloud-dataflow-java + ${beam.version} + runtime + + + + + spark-runner + + + + org.apache.beam + beam-runners-spark-3 + ${beam.version} + runtime + com.google.cloud.bigdataoss gcsio 3.1.1 - - - - - - org.apache.maven.plugins - maven-shade-plugin - - false - - - - *:* - - META-INF/*.SF - META-INF/*.DSA - META-INF/*.RSA - - - - - - - package - - shade - - + + + + + org.apache.maven.plugins + maven-shade-plugin + + false + + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + + + + + + + package + + shade + + - true - - shaded - - - - - - com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample - - - - + true + + shaded + + + + + com.google.privacy.differentialprivacy.pipelinedp4j.examples.BeamExample + + + + - + within the shaded JAR. --> com.google.common custom.guava.cache.com.google.common - - - - - - - - - - - - spark-runner-embeddeed - - - - org.apache.spark - spark-core_2.12 - ${spark.version} - - - - org.apache.spark - spark-streaming_2.12 - ${spark.version} - - - - - - + executable JAR, allowing the pipeline to run standalone. --> + + + org.apache.spark + spark-core_2.12 + ${spark.version} + + + org.apache.spark + spark-streaming_2.12 + ${spark.version} + + + + + \ No newline at end of file From c9fd92b9bb51d1b1afb0edd164f78e9b6b9032e2 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Tue, 21 Apr 2026 20:53:07 +0000 Subject: [PATCH 20/24] reformat pom --- examples/pipelinedp4j/beam/pom.xml | 29 ++++++++++++++++------------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index be7bb224..801ef952 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -1,21 +1,24 @@ - - + + + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://www.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 com.google.privacy.differentialprivacy From 6051f51ccac89607ac3bf9f6721a5f2c758f3ae5 Mon Sep 17 00:00:00 2001 From: leroyjb Date: Tue, 21 Apr 2026 20:56:18 +0000 Subject: [PATCH 21/24] reformat pom --- examples/pipelinedp4j/beam/pom.xml | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/examples/pipelinedp4j/beam/pom.xml b/examples/pipelinedp4j/beam/pom.xml index 801ef952..65e2f891 100644 --- a/examples/pipelinedp4j/beam/pom.xml +++ b/examples/pipelinedp4j/beam/pom.xml @@ -18,15 +18,18 @@ + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 + com.google.privacy.differentialprivacy pipelinedp4j 1.0-SNAPSHOT + beam jar + 2.72.0 3.5.8 @@ -35,6 +38,7 @@ UTF-8 + @@ -49,6 +53,7 @@ ${beam.version} + direct-runner @@ -65,6 +70,7 @@ + dataflow-runner @@ -77,6 +83,7 @@ + spark-runner @@ -157,6 +164,7 @@ + spark-runner-embeddeed + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 @@ -46,13 +46,14 @@ common 1.0-SNAPSHOT + org.apache.beam beam-sdks-java-core ${beam.version} - + @@ -164,7 +165,7 @@ - + spark-runner-embeddeed