diff --git a/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json b/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json new file mode 100644 index 000000000000..5abe02fc09c7 --- /dev/null +++ b/.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json @@ -0,0 +1,4 @@ +{ + "comment": "Modify this file in a trivial way to cause this test suite to run.", + "modification": 1 +} diff --git a/.github/workflows/README.md b/.github/workflows/README.md index 93fda448521b..dcc885feec9d 100644 --- a/.github/workflows/README.md +++ b/.github/workflows/README.md @@ -330,6 +330,7 @@ PostCommit Jobs run in a schedule against master branch and generally do not get | [ PostCommit Java BigQueryEarlyRollout ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_BigQueryEarlyRollout.yml) | N/A |`beam_PostCommit_Java_BigQueryEarlyRollout.json`| [![.github/workflows/beam_PostCommit_Java_BigQueryEarlyRollout.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_BigQueryEarlyRollout.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_BigQueryEarlyRollout.yml?query=event%3Aschedule) | | [ PostCommit Java Dataflow V1 ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV1.yml) | N/A |`beam_PostCommit_Java_DataflowV1.json`| [![.github/workflows/beam_PostCommit_Java_DataflowV1.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV1.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV1.yml?query=event%3Aschedule) | | [ PostCommit Java Dataflow V2 ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV2.yml) | N/A |`beam_PostCommit_Java_DataflowV2.json`| [![.github/workflows/beam_PostCommit_Java_DataflowV2.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV2.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_DataflowV2.yml?query=event%3Aschedule) | +| [ PostCommit Java Delta IO Dataflow ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml) | N/A |`beam_PostCommit_Java_Delta_IO_Dataflow.json`| [![.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml?query=event%3Aschedule) | | [ PostCommit Java Examples Dataflow ARM ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_ARM.yml) | ['8','11','17','21','25'] |`beam_PostCommit_Java_Examples_Dataflow_ARM.json`| [![.github/workflows/beam_PostCommit_Java_Examples_Dataflow_ARM.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_ARM.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_ARM.yml?query=event%3Aschedule) | | [ PostCommit Java Examples Dataflow](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow.yml) | N/A |`beam_PostCommit_Java_Examples_Dataflow.json`| [![.github/workflows/beam_PostCommit_Java_Examples_Dataflow.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow.yml?query=event%3Aschedule) | | [ PostCommit Java Examples Dataflow Java ](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_Java.yml) | ['8','17','21','25'] |`beam_PostCommit_Java_Examples_Dataflow_Java.json`| [![.github/workflows/beam_PostCommit_Java_Examples_Dataflow_Java.yml](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_Java.yml/badge.svg?event=schedule)](https://github.com/apache/beam/actions/workflows/beam_PostCommit_Java_Examples_Dataflow_Java.yml?query=event%3Aschedule) | diff --git a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml new file mode 100644 index 000000000000..604be67d90fa --- /dev/null +++ b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml @@ -0,0 +1,78 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You 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. + +name: PostCommit Java Delta IO Dataflow + +on: + schedule: + - cron: '30 4/6 * * *' + pull_request_target: + paths: [ 'release/trigger_all_tests.json', '.github/trigger_files/beam_PostCommit_Java_Delta_IO_Dataflow.json' ] + workflow_dispatch: + +# Setting explicit permissions for the action to avoid the default permissions which are `write-all` in case of pull_request_target event +permissions: + actions: write + pull-requests: write + checks: write + contents: read + deployments: read + id-token: none + issues: write + discussions: read + packages: read + pages: read + repository-projects: read + security-events: read + statuses: read + +# This allows a subsequently queued workflow run to interrupt previous runs +concurrency: + group: '${{ github.workflow }} @ ${{ github.event.pull_request.number || github.sha || github.head_ref || github.ref }}-${{ github.event.schedule || github.event.comment.id || github.event.sender.login }}' + cancel-in-progress: true + +env: + DEVELOCITY_ACCESS_KEY: ${{ secrets.DEVELOCITY_ACCESS_KEY }} + GRADLE_ENTERPRISE_CACHE_USERNAME: ${{ secrets.GE_CACHE_USERNAME }} + GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }} + +jobs: + beam_PostCommit_Java_Delta_IO_Dataflow: + if: | + github.event_name == 'pull_request_target' || + github.event_name == 'workflow_dispatch' || + (github.event_name == 'schedule' && github.repository == 'apache/beam') || + github.event.comment.body == 'Run PostCommit Java Delta IO Dataflow' + runs-on: [self-hosted, ubuntu-24.04, main] + timeout-minutes: 120 + name: ${{ matrix.job_name }} (${{ matrix.job_phrase }}) + strategy: + matrix: + job_name: ["beam_PostCommit_Java_Delta_IO_Dataflow"] + job_phrase: ["Run PostCommit Java Delta IO Dataflow"] + steps: + - uses: actions/checkout@v6 + - name: Setup repository + uses: ./.github/actions/setup-action + with: + comment_phrase: ${{ matrix.job_phrase }} + github_token: ${{ secrets.GITHUB_TOKEN }} + github_job: ${{ matrix.job_name }} (${{ matrix.job_phrase }}) + - name: Setup environment + uses: ./.github/actions/setup-environment-action + - name: Run PostCommit Java Delta IO Dataflow + uses: ./.github/actions/gradle-command-self-hosted-action + with: + gradle-command: :sdks:java:io:delta:dataflowIntegrationTest --info diff --git a/sdks/java/io/delta/build.gradle b/sdks/java/io/delta/build.gradle index 57b1cd8ad876..5ee5442ecd18 100644 --- a/sdks/java/io/delta/build.gradle +++ b/sdks/java/io/delta/build.gradle @@ -16,6 +16,8 @@ * limitations under the License. */ +import groovy.json.JsonOutput + plugins { id 'org.apache.beam.module' } applyJavaNature( automaticModuleName: 'org.apache.beam.sdk.io.delta', @@ -68,6 +70,10 @@ dependencies { testImplementation project(":sdks:java:managed") testRuntimeOnly "org.yaml:snakeyaml:2.0" testImplementation project(path: ":runners:direct-java", configuration: "shadow") + testImplementation project(":sdks:java:io:google-cloud-platform") + testImplementation project(":sdks:java:extensions:google-cloud-platform-core") + testImplementation library.java.google_cloud_storage + testRuntimeOnly project(path: ":runners:google-cloud-dataflow-java") } configurations.all { @@ -82,3 +88,35 @@ configurations.all { resolutionStrategy.force "com.google.cloud.bigdataoss:gcsio:$bigdataoss_gcs_connector_version" resolutionStrategy.force "com.google.cloud.bigdataoss:util:$bigdataoss_gcs_connector_version" } + +def gcpProject = project.findProperty('gcpProject') ?: 'apache-beam-testing' +def gcpTempLocation = project.findProperty('gcpTempLocation') ?: 'gs://apache-beam-testing-delta-lake/temp' + +// Used by Dataflow integration tests. +task dataflowIntegrationTest(type: Test) { + group = "Verification" + evaluationDependsOn(":runners:google-cloud-dataflow-java") + dependsOn ":runners:google-cloud-dataflow-java:buildAndPushDockerJavaContainer" + finalizedBy ":runners:google-cloud-dataflow-java:cleanUpDockerJavaImages" + def dockerJavaImageName = project.project(':runners:google-cloud-dataflow-java').ext.dockerJavaImageName + + def args = [ + "--runner=DataflowRunner", + "--region=us-central1", + "--project=${gcpProject}", + "--tempLocation=${gcpTempLocation}", + "--tempRoot=${gcpTempLocation}", + "--sdkContainerImage=${dockerJavaImageName}", + "--experiments=use_runner_v2,use_staged_dataflow_worker_jar" + ] + systemProperty "beamTestPipelineOptions", JsonOutput.toJson(args) + + // Disable Gradle cache: these ITs interact with live service that should always be considered "out of date" + outputs.upToDateWhen { false } + + include '**/*IT.class' + + maxParallelForks 4 + classpath = sourceSets.test.runtimeClasspath + testClassesDirs = sourceSets.test.output.classesDirs +} diff --git a/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java new file mode 100644 index 000000000000..e0d35f30faa6 --- /dev/null +++ b/sdks/java/io/delta/src/test/java/org/apache/beam/sdk/io/delta/DeltaIOIT.java @@ -0,0 +1,264 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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. + */ +package org.apache.beam.sdk.io.delta; + +import com.google.cloud.storage.Blob; +import com.google.cloud.storage.Storage; +import com.google.cloud.storage.StorageOptions; +import io.delta.kernel.DataWriteContext; +import io.delta.kernel.Operation; +import io.delta.kernel.Table; +import io.delta.kernel.Transaction; +import io.delta.kernel.TransactionBuilder; +import io.delta.kernel.TransactionCommitResult; +import io.delta.kernel.data.ColumnVector; +import io.delta.kernel.data.ColumnarBatch; +import io.delta.kernel.data.FilteredColumnarBatch; +import io.delta.kernel.defaults.engine.DefaultEngine; +import io.delta.kernel.defaults.internal.data.DefaultColumnarBatch; +import io.delta.kernel.engine.Engine; +import io.delta.kernel.types.DataType; +import io.delta.kernel.types.IntegerType; +import io.delta.kernel.types.StringType; +import io.delta.kernel.types.StructType; +import io.delta.kernel.utils.CloseableIterable; +import io.delta.kernel.utils.CloseableIterator; +import io.delta.kernel.utils.DataFileStatus; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import org.apache.beam.sdk.managed.Managed; +import org.apache.beam.sdk.schemas.Schema; +import org.apache.beam.sdk.testing.PAssert; +import org.apache.beam.sdk.testing.TestPipeline; +import org.apache.beam.sdk.values.PCollection; +import org.apache.beam.sdk.values.Row; +import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap; +import org.apache.hadoop.conf.Configuration; +import org.junit.After; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TestName; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** Integration tests for {@link DeltaIO}. */ +@RunWith(JUnit4.class) +public class DeltaIOIT { + private static final Logger LOG = LoggerFactory.getLogger(DeltaIOIT.class); + + private static final String DEFAULT_BUCKET = "apache-beam-testing-delta-lake"; + @Rule public final TestPipeline readPipeline = TestPipeline.create(); + @Rule public final TestName testName = new TestName(); + + private String bucket; + private String repoPath; + private String repoPrefix; + private Storage storage; + + private static final Schema ROW_SCHEMA = + Schema.builder().addInt32Field("id").addStringField("name").build(); + + private static final List TEST_ROWS = + IntStream.range(0, 100) + .mapToObj(i -> Row.withSchema(ROW_SCHEMA).addValues(i, "name_" + i).build()) + .collect(Collectors.toList()); + + @Before + public void setup() throws Exception { + storage = StorageOptions.newBuilder().build().getService(); + long salt = System.currentTimeMillis(); + + String tempLocation = readPipeline.getOptions().getTempLocation(); + if (tempLocation != null && tempLocation.startsWith("gs://")) { + org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath gcsPath = + org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath.fromUri(tempLocation); + bucket = gcsPath.getBucket(); + repoPrefix = gcsPath.getObject() + "/delta_io_it/" + testName.getMethodName() + "-" + salt; + } else { + bucket = DEFAULT_BUCKET; + repoPrefix = "delta_io_it/" + testName.getMethodName() + "-" + salt; + } + repoPath = "gs://" + bucket + "/" + repoPrefix; + + LOG.info("Generating Delta Lake repository at {}", repoPath); + + Configuration configuration = new Configuration(); + configuration.set("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem"); + configuration.set( + "fs.AbstractFileSystem.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS"); + configuration.set("fs.gs.auth.type", "APPLICATION_DEFAULT"); + String project = + readPipeline + .getOptions() + .as(org.apache.beam.sdk.extensions.gcp.options.GcpOptions.class) + .getProject(); + if (project != null) { + configuration.set("fs.gs.project.id", project); + } + + Engine engine = DefaultEngine.create(configuration); + Table table = Table.forPath(engine, repoPath); + + StructType deltaSchema = + new StructType().add("id", IntegerType.INTEGER).add("name", StringType.STRING); + + TransactionBuilder txnBuilder = + table.createTransactionBuilder(engine, "DeltaIOIT", Operation.CREATE_TABLE); + txnBuilder = txnBuilder.withSchema(engine, deltaSchema); + Transaction txn = txnBuilder.build(engine); + io.delta.kernel.data.Row txnState = txn.getTransactionState(engine); + + ColumnVector idVector = + new ColumnVector() { + @Override + public DataType getDataType() { + return IntegerType.INTEGER; + } + + @Override + public int getSize() { + return TEST_ROWS.size(); + } + + @Override + public void close() {} + + @Override + public boolean isNullAt(int rowId) { + return TEST_ROWS.get(rowId).getValue("id") == null; + } + + @Override + public int getInt(int rowId) { + return TEST_ROWS.get(rowId).getInt32("id"); + } + }; + + ColumnVector nameVector = + new ColumnVector() { + @Override + public DataType getDataType() { + return StringType.STRING; + } + + @Override + public int getSize() { + return TEST_ROWS.size(); + } + + @Override + public void close() {} + + @Override + public boolean isNullAt(int rowId) { + return TEST_ROWS.get(rowId).getValue("name") == null; + } + + @Override + public String getString(int rowId) { + return TEST_ROWS.get(rowId).getString("name"); + } + }; + + ColumnVector[] vectors = new ColumnVector[] {idVector, nameVector}; + ColumnarBatch columnarBatch = new DefaultColumnarBatch(TEST_ROWS.size(), deltaSchema, vectors); + FilteredColumnarBatch filteredBatch = + new FilteredColumnarBatch(columnarBatch, Optional.empty()); + + CloseableIterator data = + io.delta.kernel.internal.util.Utils.toCloseableIterator( + Collections.singletonList(filteredBatch).iterator()); + + CloseableIterator physicalData = + Transaction.transformLogicalData(engine, txnState, data, Collections.emptyMap()); + + DataWriteContext writeContext = + Transaction.getWriteContext(engine, txnState, Collections.emptyMap()); + + CloseableIterator dataFiles = + engine + .getParquetHandler() + .writeParquetFiles( + writeContext.getTargetDirectory(), + physicalData, + writeContext.getStatisticsColumns()); + + CloseableIterator dataActions = + Transaction.generateAppendActions(engine, txnState, dataFiles, writeContext); + + CloseableIterable dataActionsIterable = + CloseableIterable.inMemoryIterable(dataActions); + + TransactionCommitResult commitResult = txn.commit(engine, dataActionsIterable); + + if (commitResult.getVersion() < 0) { + throw new RuntimeException("Table creation/write failed"); + } + + LOG.info("Successfully generated Delta Lake repository"); + } + + @After + public void teardown() { + if (storage == null) { + return; + } + LOG.info("Cleaning up Delta Lake repository at {}", repoPath); + try { + Iterable blobs = + storage.list(bucket, Storage.BlobListOption.prefix(repoPrefix)).getValues(); + blobs.forEach(b -> storage.delete(b.getBlobId())); + } catch (Exception e) { + LOG.warn("Failed to clean up GCS repository at {}", repoPath, e); + } + } + + @Test + public void testReadDeltaLakeTable() { + Map hadoopConfig = new HashMap<>(); + hadoopConfig.put("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem"); + hadoopConfig.put( + "fs.AbstractFileSystem.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS"); + String project = + readPipeline + .getOptions() + .as(org.apache.beam.sdk.extensions.gcp.options.GcpOptions.class) + .getProject(); + if (project != null) { + hadoopConfig.put("fs.gs.project.id", project); + } + + PCollection output = + readPipeline + .apply( + Managed.read(Managed.DELTA_LAKE) + .withConfig(ImmutableMap.of("table", repoPath, "hadoop_config", hadoopConfig))) + .getSinglePCollection(); + + PAssert.that(output).containsInAnyOrder(TEST_ROWS); + readPipeline.run().waitUntilFinish(); + } +}