From a1ae748e9fb139a6b7514ca42a40100f3dc44eeb Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Fri, 17 Jul 2026 20:44:16 +0000 Subject: [PATCH 1/9] try fixing missing delta lake read --- .github/workflows/upload-python-package.yml | 19 ++++++++++++++++++- .../read_from_delta_lake.py | 1 + .../read_from_delta_lake_test.py | 1 + yaml/pom.xml | 18 ++++++++++++++++++ 4 files changed, 38 insertions(+), 1 deletion(-) diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index d5d688dff1..0c1d9a42a9 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -47,6 +47,15 @@ jobs: with: persist-credentials: false + - name: Set up Java + uses: actions/setup-java@v4 + with: + distribution: 'temurin' + java-version: '17' + + - name: Build Java Dependencies + run: mvn test-compile -pl yaml + - name: Set up Python uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 with: @@ -63,7 +72,15 @@ jobs: - name: Run Tests working-directory: python/src/main/python/job-builder-util-transforms - run: poetry run pytest ../../../test/python/job-builder-util-transforms/ + env: + BEAM_SERVICE_OVERRIDES: '{"sdks:java:io:expansion-service:shadowJar": "localhost:8097"}' + run: | + # Start the expansion service in the background + mvn exec:exec -pl yaml -Dexec.classpathScope=test -Dexec.executable="java" -Dexec.args="-classpath %classpath org.apache.beam.sdk.expansion.service.ExpansionService 8097 --alsoStartLoopbackWorker=true" & + # Wait 10 seconds for the service to compile and start up + sleep 10 + # Execute pytest + poetry run pytest ../../../test/python/job-builder-util-transforms/ --test-pipeline-options="--runner=FnApiRunner --environment_type=LOOPBACK" - name: Build Package working-directory: python/src/main/python/job-builder-util-transforms diff --git a/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py b/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py index d22ca896de..47306e432c 100644 --- a/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py +++ b/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py @@ -56,6 +56,7 @@ def expand(self, pbegin): expansion_service=BeamJarExpansionService( 'sdks:java:io:expansion-service:shadowJar' ), + rearrange_based_on_discovery=True, **config, ) diff --git a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py index cd30f675d0..94db327c7c 100644 --- a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py +++ b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py @@ -69,6 +69,7 @@ def test_read_from_delta_lake_fallback(self, mock_managed_read, mock_saet): mock_saet.assert_called_once_with( identifier=DELTA_LAKE_READ_URN, expansion_service=ANY, + rearrange_based_on_discovery=True, table=table, hadoop_config=hadoop_config, ) diff --git a/yaml/pom.xml b/yaml/pom.xml index 921ea00667..d4afc8cbe4 100644 --- a/yaml/pom.xml +++ b/yaml/pom.xml @@ -120,6 +120,24 @@ 1.17.0 test + + org.apache.beam + beam-sdks-java-expansion-service + ${beam.version} + test + + + org.apache.beam + beam-sdks-java-io-delta + ${beam.version} + test + + + com.github.jbellis + jamm + 0.4.0 + test + From 45a4dbd5a1c37d05b21dd18ba8b4f606438cf0c3 Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Fri, 17 Jul 2026 20:57:36 +0000 Subject: [PATCH 2/9] add local test loopback --- .github/workflows/upload-python-package.yml | 11 ++++++----- .../read_from_delta_lake_test.py | 3 +++ 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index 0c1d9a42a9..e6c1f638ff 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -73,14 +73,15 @@ jobs: - name: Run Tests working-directory: python/src/main/python/job-builder-util-transforms env: + LOCAL_TEST_LOOPBACK: '1' BEAM_SERVICE_OVERRIDES: '{"sdks:java:io:expansion-service:shadowJar": "localhost:8097"}' run: | - # Start the expansion service in the background - mvn exec:exec -pl yaml -Dexec.classpathScope=test -Dexec.executable="java" -Dexec.args="-classpath %classpath org.apache.beam.sdk.expansion.service.ExpansionService 8097 --alsoStartLoopbackWorker=true" & - # Wait 10 seconds for the service to compile and start up - sleep 10 + # Start the expansion service in the background from the repository root + (cd "$GITHUB_WORKSPACE" && mvn exec:exec -pl yaml -Dexec.classpathScope=test -Dexec.executable="java" -Dexec.args="-classpath %classpath org.apache.beam.sdk.expansion.service.ExpansionService 8097 --alsoStartLoopbackWorker=true") & + # Wait 15 seconds for the service to compile and start up + sleep 15 # Execute pytest - poetry run pytest ../../../test/python/job-builder-util-transforms/ --test-pipeline-options="--runner=FnApiRunner --environment_type=LOOPBACK" + poetry run pytest ../../../test/python/job-builder-util-transforms/ - name: Build Package working-directory: python/src/main/python/job-builder-util-transforms diff --git a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py index 94db327c7c..143f2192c9 100644 --- a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py +++ b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py @@ -4,6 +4,9 @@ from unittest.mock import ANY, MagicMock, patch import apache_beam as beam from apache_beam.testing.test_pipeline import TestPipeline + +if os.environ.get('LOCAL_TEST_LOOPBACK') == '1': + TestPipeline.pytest_test_pipeline_options = '--runner=FnApiRunner --environment_type=LOOPBACK' from apache_beam.testing.util import assert_that, equal_to from apache_beam.transforms import managed import pyarrow as pa From aceb3cd253e48ce7021f99590b27bb3dc155906d Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Fri, 17 Jul 2026 21:25:22 +0000 Subject: [PATCH 3/9] add todo to remove java code build --- .github/workflows/upload-python-package.yml | 3 +++ yaml/pom.xml | 1 + 2 files changed, 4 insertions(+) diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index e6c1f638ff..181c0299a9 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -47,6 +47,9 @@ jobs: with: persist-credentials: false + # TODO(BEAM-2.76.0): Remove Java setup and background expansion service steps once Beam is upgraded to 2.76.0, + # as the official 2.76.0 expansion service shadowJar will natively package the Delta Lake IO classes. + # See PR#4038 for further details on code removal. - name: Set up Java uses: actions/setup-java@v4 with: diff --git a/yaml/pom.xml b/yaml/pom.xml index d4afc8cbe4..875dd7e29a 100644 --- a/yaml/pom.xml +++ b/yaml/pom.xml @@ -120,6 +120,7 @@ 1.17.0 test + org.apache.beam beam-sdks-java-expansion-service From 728d412b3152fae7e4e913bdbf6b1473173df26a Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Fri, 17 Jul 2026 22:36:53 +0000 Subject: [PATCH 4/9] fix gemini comment --- .../read_from_delta_lake_test.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py index 143f2192c9..49b9458849 100644 --- a/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py +++ b/python/src/test/python/job-builder-util-transforms/read_from_delta_lake_test.py @@ -4,15 +4,15 @@ from unittest.mock import ANY, MagicMock, patch import apache_beam as beam from apache_beam.testing.test_pipeline import TestPipeline - -if os.environ.get('LOCAL_TEST_LOOPBACK') == '1': - TestPipeline.pytest_test_pipeline_options = '--runner=FnApiRunner --environment_type=LOOPBACK' from apache_beam.testing.util import assert_that, equal_to from apache_beam.transforms import managed import pyarrow as pa import pyarrow.parquet as pq from read_from_delta_lake import DELTA_LAKE_READ_URN, ReadFromDeltaLake +if os.environ.get('LOCAL_TEST_LOOPBACK') == '1': + TestPipeline.pytest_test_pipeline_options = '--runner=FnApiRunner --environment_type=LOOPBACK' + class ReadFromDeltaLakeTest(unittest.TestCase): From f10f266f2a6fd906f2b77f97e871b6cb202efe72 Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Mon, 20 Jul 2026 14:40:36 +0000 Subject: [PATCH 5/9] build expansion jar that includes delta lake --- .github/workflows/upload-python-package.yml | 3 +- yaml/pom.xml | 38 +++++++++++++++++++-- 2 files changed, 37 insertions(+), 4 deletions(-) diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index 181c0299a9..a320d02832 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -57,7 +57,7 @@ jobs: java-version: '17' - name: Build Java Dependencies - run: mvn test-compile -pl yaml + run: mvn package -pl yaml -am -DskipTests - name: Set up Python uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6.3.0 @@ -101,6 +101,7 @@ jobs: GCS_PACKAGE_PATH="$BUCKET_PATH/$DATE/job_builder_util_transforms-$VERSION.tar.gz" if [ "$ONLY_UPLOAD_LISTING" != "true" ]; then gcloud storage cp python/src/main/python/job-builder-util-transforms/dist/job_builder_util_transforms-*.tar.gz "$GCS_PACKAGE_PATH" + gcloud storage cp yaml/target/template-yaml-1.0-SNAPSHOT-shaded.jar "$BUCKET_PATH/$DATE/expansion-service-custom-$VERSION.jar" fi PACKAGE_PATH_NO_GS="${GCS_PACKAGE_PATH#gs://}" sed "s||$PACKAGE_PATH_NO_GS|g" python/src/main/python/job-builder-util-transforms/provider_listing_template.yaml > "python/src/main/python/job-builder-util-transforms/yaml_provider_listing_$VERSION.yaml" diff --git a/yaml/pom.xml b/yaml/pom.xml index 875dd7e29a..8762e0b433 100644 --- a/yaml/pom.xml +++ b/yaml/pom.xml @@ -125,21 +125,22 @@ org.apache.beam beam-sdks-java-expansion-service ${beam.version} - test + runtime org.apache.beam beam-sdks-java-io-delta ${beam.version} - test + runtime com.github.jbellis jamm 0.4.0 - test + runtime + @@ -156,6 +157,37 @@ + + + org.apache.maven.plugins + maven-shade-plugin + 3.5.1 + + + package + + shade + + + true + shaded + + + + + + *:* + + META-INF/*.SF + META-INF/*.DSA + META-INF/*.RSA + + + + + + + From 41378390c113270571ea83b3f55044c4b6c07b4b Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Mon, 20 Jul 2026 15:17:02 +0000 Subject: [PATCH 6/9] update running of java jar --- .github/workflows/upload-python-package.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index a320d02832..b778b0f140 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -79,9 +79,9 @@ jobs: LOCAL_TEST_LOOPBACK: '1' BEAM_SERVICE_OVERRIDES: '{"sdks:java:io:expansion-service:shadowJar": "localhost:8097"}' run: | - # Start the expansion service in the background from the repository root - (cd "$GITHUB_WORKSPACE" && mvn exec:exec -pl yaml -Dexec.classpathScope=test -Dexec.executable="java" -Dexec.args="-classpath %classpath org.apache.beam.sdk.expansion.service.ExpansionService 8097 --alsoStartLoopbackWorker=true") & - # Wait 15 seconds for the service to compile and start up + # Start the expansion service in the background from the shaded jar + java -cp "$GITHUB_WORKSPACE/yaml/target/template-yaml-1.0-SNAPSHOT-shaded.jar" org.apache.beam.sdk.expansion.service.ExpansionService 8097 --alsoStartLoopbackWorker=true & + # Wait 15 seconds for the service to start up sleep 15 # Execute pytest poetry run pytest ../../../test/python/job-builder-util-transforms/ From 072f421028f92f3b33384e9f03467cd5cbaa4a5f Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Mon, 20 Jul 2026 15:23:14 +0000 Subject: [PATCH 7/9] add more dependencies at runtime --- yaml/pom.xml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/yaml/pom.xml b/yaml/pom.xml index 8762e0b433..74c0de580e 100644 --- a/yaml/pom.xml +++ b/yaml/pom.xml @@ -94,7 +94,7 @@ org.apache.iceberg iceberg-gcp 1.10.1 - test + runtime org.json @@ -106,19 +106,19 @@ org.apache.avro avro ${avro.version} - test + runtime org.apache.parquet parquet-hadoop 1.17.0 - test + runtime org.apache.parquet parquet-avro 1.17.0 - test + runtime From 5be2a6df30d240babe614ceb6f7b2b92c183d4b4 Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Mon, 20 Jul 2026 17:14:52 +0000 Subject: [PATCH 8/9] point delta lake to the created expansion service jar --- .../job-builder-util-transforms/pyproject.toml | 2 +- .../read_from_delta_lake.py | 17 ++++++++++++++--- 2 files changed, 15 insertions(+), 4 deletions(-) diff --git a/python/src/main/python/job-builder-util-transforms/pyproject.toml b/python/src/main/python/job-builder-util-transforms/pyproject.toml index 48ee860e6f..bc61f4aa68 100644 --- a/python/src/main/python/job-builder-util-transforms/pyproject.toml +++ b/python/src/main/python/job-builder-util-transforms/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "job-builder-util-transforms" -version = "0.1.0" +version = "0.2.0" description = "Utility transforms for job builder" authors = ["Google Cloud Platform"] packages = [ diff --git a/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py b/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py index 47306e432c..7fcdfa974b 100644 --- a/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py +++ b/python/src/main/python/job-builder-util-transforms/read_from_delta_lake.py @@ -1,9 +1,11 @@ """Module containing transforms to read data from Delta Lake tables.""" from typing import Mapping, Optional +from apache_beam.options.pipeline_options import CrossLanguageOptions from apache_beam.transforms import PTransform from apache_beam.transforms import managed from apache_beam.transforms.external import BeamJarExpansionService +from apache_beam.transforms.external import JavaJarExpansionService from apache_beam.transforms.external import SchemaAwareExternalTransform DELTA_LAKE_READ_URN = "beam:schematransform:org.apache.beam:delta_lake_read:v1" @@ -51,11 +53,20 @@ def expand(self, pbegin): ): return pbegin | managed.Read(delta_source, config=config) else: + options = pbegin.pipeline.options + beam_services = options.view_as(CrossLanguageOptions).beam_services or {} + if 'sdks:java:io:expansion-service:shadowJar' in beam_services: + expansion_service = BeamJarExpansionService( + 'sdks:java:io:expansion-service:shadowJar' + ) + else: + expansion_service = JavaJarExpansionService( + 'https://storage.googleapis.com/dataflow-templates/extra-python-packages/2026-07-20/expansion-service-custom-0.2.0.jar' + ) + return pbegin | SchemaAwareExternalTransform( identifier=DELTA_LAKE_READ_URN, - expansion_service=BeamJarExpansionService( - 'sdks:java:io:expansion-service:shadowJar' - ), + expansion_service=expansion_service, rearrange_based_on_discovery=True, **config, ) From 78af2cfbd6624f420335734dfeaeba375d8fb271 Mon Sep 17 00:00:00 2001 From: Derrick Williams Date: Mon, 20 Jul 2026 17:38:46 +0000 Subject: [PATCH 9/9] revert pyproject.toml and save for next PR --- .../src/main/python/job-builder-util-transforms/pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/src/main/python/job-builder-util-transforms/pyproject.toml b/python/src/main/python/job-builder-util-transforms/pyproject.toml index bc61f4aa68..48ee860e6f 100644 --- a/python/src/main/python/job-builder-util-transforms/pyproject.toml +++ b/python/src/main/python/job-builder-util-transforms/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "job-builder-util-transforms" -version = "0.2.0" +version = "0.1.0" description = "Utility transforms for job builder" authors = ["Google Cloud Platform"] packages = [