diff --git a/.github/workflows/upload-python-package.yml b/.github/workflows/upload-python-package.yml index c4ded12686..adbbcb0522 100644 --- a/.github/workflows/upload-python-package.yml +++ b/.github/workflows/upload-python-package.yml @@ -47,6 +47,18 @@ 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: + distribution: 'temurin' + java-version: '17' + + - name: Build Java Dependencies + run: mvn package -pl yaml -am -DskipTests + - name: Set up Python uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0 with: @@ -63,7 +75,16 @@ 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: + 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 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/ - name: Build Package working-directory: python/src/main/python/job-builder-util-transforms @@ -80,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/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..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,21 @@ 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, ) 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..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 @@ -10,6 +10,9 @@ 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): @@ -69,6 +72,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..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,20 +106,40 @@ 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 + + + + org.apache.beam + beam-sdks-java-expansion-service + ${beam.version} + runtime + + org.apache.beam + beam-sdks-java-io-delta + ${beam.version} + runtime + + + com.github.jbellis + jamm + 0.4.0 + runtime + + @@ -137,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 + + + + + + +