diff --git a/.asf.yaml b/.asf.yaml index b650499326d9..7adaf351c8f5 100644 --- a/.asf.yaml +++ b/.asf.yaml @@ -51,6 +51,7 @@ github: protected_branches: master: {} + release-2.75: {} release-2.74.0-postrelease: {} release-2.74: {} release-2.73.0-postrelease: {} diff --git a/.github/trigger_files/beam_PostCommit_Go_VR_Flink.json b/.github/trigger_files/beam_PostCommit_Go_VR_Flink.json index 939c43396fda..504c32974c7c 100644 --- a/.github/trigger_files/beam_PostCommit_Go_VR_Flink.json +++ b/.github/trigger_files/beam_PostCommit_Go_VR_Flink.json @@ -1,6 +1,6 @@ { "comment": "Modify this file in a trivial way to cause this test suite to run", - "modification": 2, + "modification": 3, "https://github.com/apache/beam/pull/32440": "testing datastream optimizations", "pr": "37640" } diff --git a/.github/trigger_files/beam_PostCommit_XVR_Flink.json b/.github/trigger_files/beam_PostCommit_XVR_Flink.json index a926b3314ed6..7dcd6398db10 100644 --- a/.github/trigger_files/beam_PostCommit_XVR_Flink.json +++ b/.github/trigger_files/beam_PostCommit_XVR_Flink.json @@ -1,4 +1,4 @@ { - "modification": 2, + "modification": 3, "trigger-2026-04-04": "portable_runner expand_sdf opt-in" } diff --git a/.github/trigger_files/beam_PostCommit_XVR_Samza.json b/.github/trigger_files/beam_PostCommit_XVR_Samza.json deleted file mode 100644 index 2bf3f556083b..000000000000 --- a/.github/trigger_files/beam_PostCommit_XVR_Samza.json +++ /dev/null @@ -1 +0,0 @@ -{"modification": 2} \ No newline at end of file diff --git a/.github/workflows/cut_release_branch.yml b/.github/workflows/cut_release_branch.yml index ee571874ce13..8b591524ba55 100644 --- a/.github/workflows/cut_release_branch.yml +++ b/.github/workflows/cut_release_branch.yml @@ -116,7 +116,10 @@ jobs: git config user.name $GITHUB_ACTOR git config user.email actions@"$RUNNER_NAME".local - name: Install xmllint - run: sudo apt-get install -y libxml2-utils + run: | + sudo apt-get clean + sudo apt-get update + sudo apt-get install -y --no-install-recommends libxml2-utils - name: Update .asf.yaml to protect new release branch from force push run: | sed -i -e "s/master: {}/master: {}\n release-${RELEASE}: {}/g" .asf.yaml diff --git a/CHANGES.md b/CHANGES.md index ac571c2f55ae..6d7bad8c97f9 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -75,8 +75,8 @@ * (Java) Enabled state tag encoding v2 by default for new Dataflow Streaming Engine jobs. It can be disabled by passing `--experiments=disable_streaming_engine_state_tag_encoding_v2` or `--updateCompatibilityVersion=2.74.0` pipeline option. Note that the tag encoding version cannot change during a job update. Jobs using tag encoding v2 (enabled by default for new jobs on 2.75.0+) cannot be downgraded to Beam versions prior to 2.73.0, as only versions 2.73.0 and later support tag encoding v2. ([#38705](https://github.com/apache/beam/issues/38705)). * (Python) Added instrumentation to support off-the-shelf profiling agents when launching Python SDK Harness ([#38853](https://github.com/apache/beam/issues/38853)). * (Java) Added support to the FnApi Data stream protocol allowing runners to isolate bundles slowly processing input from other bundles. ([#39001](https://github.com/apache/beam/issues/39001)). -* (Java) Flink 2.1 support added ([#38947](https://github.com/apache/beam/issues/38947)). -* (Java) Flink 2.2 support added ([#38978](https://github.com/apache/beam/issues/38978)). +* (Java) Flink 2.1 and 2.2 support is added ([#38947](https://github.com/apache/beam/issues/38947)) ([#38978](https://github.com/apache/beam/issues/38978)); Flink 1.17 and 1.18 support is dropped. +* (Python) MqttIO is now supported in Python via cross-language ([#21060](https://github.com/apache/beam/issues/21060)). ## Breaking Changes diff --git a/gradle.properties b/gradle.properties index bc84397c63ce..1200e0a9f7f5 100644 --- a/gradle.properties +++ b/gradle.properties @@ -30,8 +30,8 @@ signing.gnupg.useLegacyGpg=true # buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy. # To build a custom Beam version make sure you change it in both places, see # https://github.com/apache/beam/issues/21302. -version=2.75.0-SNAPSHOT -sdk_version=2.75.0.dev +version=2.76.0-SNAPSHOT +sdk_version=2.76.0.dev javaVersion=11 diff --git a/runners/flink/job-server/flink_job_server.gradle b/runners/flink/job-server/flink_job_server.gradle index 335439cbda91..a6abca5e8586 100644 --- a/runners/flink/job-server/flink_job_server.gradle +++ b/runners/flink/job-server/flink_job_server.gradle @@ -248,10 +248,26 @@ def setupTask = project.tasks.register("flinkJobServerSetup", Exec) { def flinkJobServerJar = shadowJar.archivePath def flinkDir = project.project(":runners:flink").projectDir def additionalArgs = "" - if (project.hasProperty('flinkConfDir')) + + if (project.hasProperty('flinkConfDir')) { additionalArgs += " --flink-conf-dir=${project.property('flinkConfDir')}" - else + } + else if (isFlink2) { + def flinkConfDir = "$flinkDir/2.0/src/test/resources" + additionalArgs += "--flink-conf-dir=${project.buildDir}/flink-conf" + + doFirst { + copy { + from "$flinkDir/2.0/src/test/resources/flink-test-config.yaml" + into "${project.buildDir}/flink-conf" + + // Rename the file during the copy process + rename 'flink-test-config.yaml', 'config.yaml' + } + } + } else { additionalArgs += "--flink-conf-dir=$flinkDir/src/test/resources" + } executable 'sh' args '-c', "$pythonDir/scripts/run_job_server.sh stop --group_id ${project.name} && $pythonDir/scripts/run_job_server.sh start --group_id ${project.name} --job_port ${jobPort} --artifact_port ${artifactPort} --job_server_jar ${flinkJobServerJar} --additional_args \"${additionalArgs}\"" diff --git a/runners/google-cloud-dataflow-java/build.gradle b/runners/google-cloud-dataflow-java/build.gradle index c569c03b956e..3d1f8d78777c 100644 --- a/runners/google-cloud-dataflow-java/build.gradle +++ b/runners/google-cloud-dataflow-java/build.gradle @@ -52,8 +52,8 @@ evaluationDependsOn(":sdks:java:container:java11") ext.dataflowLegacyEnvironmentMajorVersion = '8' ext.dataflowFnapiEnvironmentMajorVersion = '8' -ext.dataflowLegacyContainerVersion = 'beam-master-20260601' -ext.dataflowFnapiContainerVersion = 'beam-master-20260601' +ext.dataflowLegacyContainerVersion = 'beam-master-20260624' +ext.dataflowFnapiContainerVersion = 'beam-master-20260624' ext.dataflowContainerBaseRepository = 'gcr.io/cloud-dataflow/v1beta3' processResources { @@ -486,6 +486,9 @@ def validatesRunnerStreamingConfig = [ 'org.apache.beam.sdk.transforms.ParDoLifecycleTest.testTeardownCalledAfterExceptionInSetupStateful', 'org.apache.beam.sdk.transforms.ParDoLifecycleTest.testTeardownCalledAfterExceptionInStartBundle', 'org.apache.beam.sdk.transforms.ParDoLifecycleTest.testTeardownCalledAfterExceptionInStartBundleStateful', + // Bundle finalizer tests are flaky https://github.com/apache/beam/issues/38710 + 'org.apache.beam.sdk.transforms.SplittableDoFnTest.testBundleFinalizationOccursOnBoundedSplittableDoFn', + 'org.apache.beam.sdk.transforms.SplittableDoFnTest.testBundleFinalizationOccursOnUnboundedSplittableDoFn', ] ] diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 9f063d393703..180dda153bb6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -1061,6 +1061,12 @@ public static void main(String[] args) throws Exception { } }); LOG.info("Enabled Open Telemetry with properties: {}", openTelemetryProperties); + } else { + // turn off auth extension so it doesn't interfere if user is configuring otel e.g. via + // JvmInitializer. + if (System.getProperty("google.otel.auth.target.signals") == null) { + System.setProperty("google.otel.auth.target.signals", "none"); + } } LOG.debug("Creating StreamingDataflowWorker from options: {}", options); diff --git a/scripts/beam-sql.sh b/scripts/beam-sql.sh index 9df48f6c18d4..906812eb1597 100755 --- a/scripts/beam-sql.sh +++ b/scripts/beam-sql.sh @@ -22,7 +22,7 @@ set -e # Exit immediately if a command exits with a non-zero status. # --- Configuration --- -DEFAULT_BEAM_VERSION="2.75.0" +DEFAULT_BEAM_VERSION="2.76.0" MAIN_CLASS="org.apache.beam.sdk.extensions.sql.jdbc.BeamSqlLine" # Directory to store cached executable JAR files CACHE_DIR="${HOME}/.beam/cache" diff --git a/sdks/go/pkg/beam/core/core.go b/sdks/go/pkg/beam/core/core.go index e3a8900b1e77..c9c26a7311cd 100644 --- a/sdks/go/pkg/beam/core/core.go +++ b/sdks/go/pkg/beam/core/core.go @@ -27,7 +27,7 @@ const ( // SdkName is the human readable name of the SDK for UserAgents. SdkName = "Apache Beam SDK for Go" // SdkVersion is the current version of the SDK. - SdkVersion = "2.75.0.dev" + SdkVersion = "2.76.0.dev" // DefaultDockerImage represents the associated image for this release. DefaultDockerImage = "apache/beam_go_sdk:" + SdkVersion diff --git a/sdks/go/test/build.gradle b/sdks/go/test/build.gradle index 3437ba12f2c7..8fcb09166146 100644 --- a/sdks/go/test/build.gradle +++ b/sdks/go/test/build.gradle @@ -92,7 +92,7 @@ task flinkValidatesRunner { doFirst { // Copy Flink conf file copy { - from "${project.rootDir}/runners/flink/${flinkVersion}/src/test/resources/flink-test-config.yaml" + from "${project.rootDir}/runners/flink/2.0/src/test/resources/flink-test-config.yaml" into "${project.buildDir}/flink-conf" // Rename the file during the copy process diff --git a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java index 60e83251f147..7bda0e18cad8 100644 --- a/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java +++ b/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java @@ -315,6 +315,12 @@ public static void main( } }); LOG.info("Enabled Open Telemetry with properties: {}", openTelemetryProperties); + } else { + // turn off auth extension so it doesn't interfere if user is configuring otel e.g. via + // JvmInitializer. + if (System.getProperty("google.otel.auth.target.signals") == null) { + System.setProperty("google.otel.auth.target.signals", "none"); + } } EnumMap< BeamFnApi.InstructionRequest.RequestCase, diff --git a/sdks/python/apache_beam/io/gcp/bigquery_file_loads.py b/sdks/python/apache_beam/io/gcp/bigquery_file_loads.py index 4e45d0324ee2..4ef6c392254b 100644 --- a/sdks/python/apache_beam/io/gcp/bigquery_file_loads.py +++ b/sdks/python/apache_beam/io/gcp/bigquery_file_loads.py @@ -491,6 +491,8 @@ class TriggerCopyJobs(beam.DoFn): """ TRIGGER_DELETE_TEMP_TABLES = 'TriggerDeleteTempTables' + # https://docs.cloud.google.com/bigquery/quotas#copy_jobs + MAX_SOURCES_PER_COPY_JOB = 1200 def __init__( self, @@ -528,96 +530,90 @@ def process( self, element_list, job_name_prefix=None, unused_schema_mod_jobs=None): if isinstance(element_list, tuple): # Allow this for streaming update compatibility while fixing BEAM-24535. - self.process_one(element_list, job_name_prefix) - else: - for element in element_list: - self.process_one(element, job_name_prefix) + element_list = [element_list] - def process_one(self, element, job_name_prefix): - destination, job_reference = element + if not element_list: + return - copy_to_reference = bigquery_tools.parse_table_reference(destination) + first_destination = element_list[0][0] + copy_to_reference = bigquery_tools.parse_table_reference(first_destination) if copy_to_reference.projectId is None: copy_to_reference.projectId = vp.RuntimeValueProvider.get_value( 'project', str, '') or self.project - copy_from_reference = bigquery_tools.parse_table_reference(destination) - copy_from_reference.tableId = job_reference.jobId - if copy_from_reference.projectId is None: - copy_from_reference.projectId = vp.RuntimeValueProvider.get_value( - 'project', str, '') or self.project - - _LOGGER.info( - "Triggering copy job from %s to %s", - copy_from_reference, - copy_to_reference) + copy_from_references = [] + for destination, job_reference in element_list: + copy_from_reference = bigquery_tools.parse_table_reference(destination) + copy_from_reference.tableId = job_reference.jobId + if copy_from_reference.projectId is None: + copy_from_reference.projectId = vp.RuntimeValueProvider.get_value( + 'project', str, '') or self.project + copy_from_references.append(copy_from_reference) - wait_for_job, write_disposition = ( - self._determine_write_disposition(copy_to_reference)) + full_table_ref = bigquery_tools.get_hashable_destination(copy_to_reference) - if not self.bq_io_metadata: - self.bq_io_metadata = create_bigquery_io_metadata(self._step_name) + is_first_time = full_table_ref not in self._observed_tables + if is_first_time: + self._observed_tables.add(full_table_ref) + if self.bq_io_metadata: + Lineage.sinks().add( + 'bigquery', + copy_to_reference.projectId, + copy_to_reference.datasetId, + copy_to_reference.tableId) + + # Split into chunks of MAX_SOURCES_PER_COPY_JOB + chunks = [ + copy_from_references[i:i + self.MAX_SOURCES_PER_COPY_JOB] + for i in range( + 0, len(copy_from_references), self.MAX_SOURCES_PER_COPY_JOB) + ] + + copy_job_name_base = '%s_%s' % ( + job_name_prefix, + _bq_uuid(bigquery_tools.get_hashable_destination(copy_to_reference))) project_id = ( copy_to_reference.projectId if self.load_job_project_id is None else self.load_job_project_id) - copy_job_name = '%s_%s' % ( - job_name_prefix, - _bq_uuid( - '%s:%s.%s' % ( - copy_from_reference.projectId, - copy_from_reference.datasetId, - copy_from_reference.tableId))) - job_reference = self.bq_wrapper._insert_copy_job( - project_id, - copy_job_name, - copy_from_reference, - copy_to_reference, - create_disposition=self.create_disposition, - write_disposition=write_disposition, - job_labels=self.bq_io_metadata.add_additional_bq_job_labels()) - - if wait_for_job: - self.bq_wrapper.wait_for_bq_job(job_reference, sleep_duration_sec=10) - self.pending_jobs.append( - GlobalWindows.windowed_value((destination, job_reference))) - def _determine_write_disposition(self, copy_to_reference) -> tuple[bool, str]: - """ - Determines the write disposition for a BigQuery copy job, - based on destination. - - When the write_disposition for a job is WRITE_TRUNCATE, multiple copy jobs - to the same destination can interfere with each other, truncate data, and - write to the BigQuery table repeatedly. To prevent this, the first copy job - runs with the user's specified write_disposition, but subsequent jobs must - always use WRITE_APPEND. This ensures that subsequent copy jobs do not - clear out data appended by previous jobs. - - Args: - copy_to_reference: The reference to the destination table. - - Returns: - A tuple containing a boolean indicating whether to wait for the job to - complete and the write disposition to use for the job. - """ - full_table_ref = '%s:%s.%s' % ( - copy_to_reference.projectId, - copy_to_reference.datasetId, - copy_to_reference.tableId) - if full_table_ref not in self._observed_tables: - write_disposition = self.write_disposition - wait_for_job = True - self._observed_tables.add(full_table_ref) - Lineage.sinks().add( - 'bigquery', - copy_to_reference.projectId, - copy_to_reference.datasetId, - copy_to_reference.tableId) - else: - wait_for_job = False - write_disposition = 'WRITE_APPEND' - return wait_for_job, write_disposition + for i, chunk in enumerate(chunks): + if i == 0 and is_first_time: + write_disposition = self.write_disposition + # Wait inline only if we have multiple chunks and write disposition is WRITE_TRUNCATE or WRITE_EMPTY. + # This ensures the first chunk initializes the table, and subsequent chunks (WRITE_APPEND) append to it. + wait_for_job = ( + self.write_disposition in ('WRITE_TRUNCATE', 'WRITE_EMPTY') and + len(chunks) > 1) + else: + write_disposition = 'WRITE_APPEND' + wait_for_job = False + + chunk_job_name = copy_job_name_base + if len(chunks) > 1: + chunk_job_name = f"{copy_job_name_base}_{i}" + + _LOGGER.info( + "Triggering copy job %s from %s to %s (write_disposition: %s)", + chunk_job_name, [str(r) for r in chunk], + copy_to_reference, + write_disposition) + + job_reference = self.bq_wrapper._insert_copy_job( + project_id, + chunk_job_name, + chunk, + copy_to_reference, + create_disposition=self.create_disposition, + write_disposition=write_disposition, + job_labels=self.bq_io_metadata.add_additional_bq_job_labels() + if self.bq_io_metadata else None) + + if wait_for_job: + self.bq_wrapper.wait_for_bq_job(job_reference, sleep_duration_sec=10) + + self.pending_jobs.append( + GlobalWindows.windowed_value((first_destination, job_reference))) def finish_bundle(self): for windowed_value in self.pending_jobs: @@ -744,7 +740,7 @@ def process( else: try: schema = bigquery_tools.table_schema_to_dict( - bigquery_tools.BigQueryWrapper().get_table( + self.bq_wrapper.get_table( project_id=table_reference.projectId, dataset_id=table_reference.datasetId, table_id=table_reference.tableId).schema) @@ -855,7 +851,8 @@ def process(self, element): if latest_partition.can_accept(file_size): latest_partition.add(file_path, file_size) else: - partitions.append(latest_partition.files) + if latest_partition.files: + partitions.append(latest_partition.files) latest_partition = PartitionFiles.Partition( self.max_partition_size, self.max_files_per_partition) latest_partition.add(file_path, file_size) @@ -1181,12 +1178,13 @@ def _load_data( # the truncation happens only once. See # https://github.com/apache/beam/issues/24535. finished_temp_tables_load_job_ids_list_pc = ( - finished_temp_tables_load_job_ids_pc | beam.MapTuple( + finished_temp_tables_load_job_ids_pc + | beam.MapTuple( lambda destination, job_reference: ( - bigquery_tools.parse_table_reference(destination).tableId, + bigquery_tools.get_hashable_destination(destination), (destination, job_reference))) | beam.GroupByKey() - | beam.MapTuple(lambda tableId, batch: list(batch))) + | beam.MapTuple(lambda dest, batch: list(batch))) else: # Loads can happen in parallel. finished_temp_tables_load_job_ids_list_pc = ( diff --git a/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py b/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py index 191719e6a208..47c1ce5ea1bb 100644 --- a/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py +++ b/sdks/python/apache_beam/io/gcp/bigquery_file_loads_test.py @@ -924,69 +924,180 @@ def dynamic_destination_resolver(element, *side_inputs): write_disposition=BigQueryDisposition.WRITE_TRUNCATE)) from apache_beam.io.gcp.internal.clients.bigquery import TableReference - mock_insert_copy_job.assert_has_calls( - [ - call( - 'project1', - mock.ANY, + mock_insert_copy_job.assert_has_calls([ + call( + 'project1', + mock.ANY, + [ TableReference( datasetId='dataset1', projectId='project1', tableId='job_name1'), - TableReference( - datasetId='dataset1', - projectId='project1', - tableId='table1'), - create_disposition=None, - write_disposition='WRITE_TRUNCATE', - job_labels={'step_name': 'bigquerybatchfileloads'}), - call( - 'project1', - mock.ANY, TableReference( datasetId='dataset1', projectId='project1', tableId='job_name1'), - TableReference( - datasetId='dataset1', - projectId='project1', - tableId='table1'), - create_disposition=None, - write_disposition='WRITE_APPEND', - job_labels={'step_name': 'bigquerybatchfileloads'}), - call( - 'project1', - mock.ANY, + ], + TableReference( + datasetId='dataset1', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_TRUNCATE', + job_labels={'step_name': 'bigquerybatchfileloads'}), + call( + 'project1', + mock.ANY, + [ TableReference( datasetId='dataset2', projectId='project1', tableId='job_name1'), - TableReference( - datasetId='dataset2', - projectId='project1', - tableId='table1'), - create_disposition=None, - # Previously this was `WRITE_APPEND`. - write_disposition='WRITE_TRUNCATE', - job_labels={'step_name': 'bigquerybatchfileloads'}), - call( - 'project1', - mock.ANY, + ], + TableReference( + datasetId='dataset2', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_TRUNCATE', + job_labels={'step_name': 'bigquerybatchfileloads'}), + call( + 'project1', + mock.ANY, + [ TableReference( datasetId='dataset3', projectId='project1', tableId='job_name1'), - TableReference( - datasetId='dataset3', - projectId='project1', - tableId='table1'), - create_disposition=None, - # Previously this was `WRITE_APPEND`. - write_disposition='WRITE_TRUNCATE', - job_labels={'step_name': 'bigquerybatchfileloads'}), - ], - any_order=True) - self.assertEqual(4, mock_insert_copy_job.call_count) + ], + TableReference( + datasetId='dataset3', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_TRUNCATE', + job_labels={'step_name': 'bigquerybatchfileloads'}), + ], + any_order=True) + self.assertEqual(3, mock_insert_copy_job.call_count) + + @mock.patch( + 'apache_beam.io.gcp.bigquery_tools.BigQueryWrapper.wait_for_bq_job') + @mock.patch( + 'apache_beam.io.gcp.bigquery_tools.BigQueryWrapper._insert_copy_job') + def test_copy_jobs_splitting( + self, mock_insert_copy_job, mock_wait_for_bq_job): + destination = 'project1:dataset1.table1' + + from apache_beam.io.gcp.bigquery_file_loads import TriggerCopyJobs + original_max_sources = TriggerCopyJobs.MAX_SOURCES_PER_COPY_JOB + TriggerCopyJobs.MAX_SOURCES_PER_COPY_JOB = 2 + + try: + job_reference = bigquery_api.JobReference() + job_reference.projectId = 'project1' + job_reference.jobId = 'job_name1' + result_job = mock.Mock() + result_job.jobReference = job_reference + + mock_job = mock.Mock() + mock_job.status.state = 'DONE' + mock_job.status.errorResult = None + mock_job.jobReference = job_reference + + bq_client = mock.Mock() + bq_client.jobs.Get.return_value = mock_job + bq_client.jobs.Insert.return_value = result_job + bq_client.tables.Delete.return_value = None + mock_insert_copy_job.return_value = job_reference + temp_dir = self._new_tempdir() + + with TestPipeline('FnApiRunner') as p: + _ = ( + p + | beam.Create([ + { + 'name': 'a' + }, + { + 'name': 'b' + }, + { + 'name': 'c' + }, + { + 'name': 'd' + }, + { + 'name': 'e' + }, + ], + reshuffle=False) + | bqfl.BigQueryBatchFileLoads( + destination, + custom_gcs_temp_location=temp_dir, + test_client=bq_client, + validate=False, + temp_file_format=bigquery_tools.FileFormat.JSON, + max_file_size=10, + max_partition_size=10, + max_files_per_partition=1, + write_disposition=BigQueryDisposition.WRITE_TRUNCATE)) + + self.assertEqual(3, mock_insert_copy_job.call_count) + + from apache_beam.io.gcp.internal.clients.bigquery import TableReference + expected_calls = [ + call( + 'project1', + mock.ANY, + [ + TableReference( + datasetId='dataset1', + projectId='project1', + tableId='job_name1'), + TableReference( + datasetId='dataset1', + projectId='project1', + tableId='job_name1'), + ], + TableReference( + datasetId='dataset1', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_TRUNCATE', + job_labels=mock.ANY), + call( + 'project1', + mock.ANY, + [ + TableReference( + datasetId='dataset1', + projectId='project1', + tableId='job_name1'), + TableReference( + datasetId='dataset1', + projectId='project1', + tableId='job_name1'), + ], + TableReference( + datasetId='dataset1', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_APPEND', + job_labels=mock.ANY), + call( + 'project1', + mock.ANY, + [ + TableReference( + datasetId='dataset1', + projectId='project1', + tableId='job_name1'), + ], + TableReference( + datasetId='dataset1', projectId='project1', tableId='table1'), + create_disposition=None, + write_disposition='WRITE_APPEND', + job_labels=mock.ANY), + ] + mock_insert_copy_job.assert_has_calls(expected_calls, any_order=True) + self.assertEqual(9, mock_wait_for_bq_job.call_count) + + finally: + TriggerCopyJobs.MAX_SOURCES_PER_COPY_JOB = original_max_sources @parameterized.expand([ param(is_streaming=False, with_auto_sharding=False, compat_version=None), diff --git a/sdks/python/apache_beam/io/gcp/bigquery_tools.py b/sdks/python/apache_beam/io/gcp/bigquery_tools.py index 8dd58cd55a01..491b7a39b0b7 100644 --- a/sdks/python/apache_beam/io/gcp/bigquery_tools.py +++ b/sdks/python/apache_beam/io/gcp/bigquery_tools.py @@ -506,16 +506,22 @@ def _insert_copy_job( reference = bigquery.JobReference() reference.jobId = job_id reference.projectId = project_id + + copy_config = bigquery.JobConfigurationTableCopy( + destinationTable=to_table_reference, + createDisposition=create_disposition, + writeDisposition=write_disposition, + ) + if isinstance(from_table_reference, list): + copy_config.sourceTables = from_table_reference + else: + copy_config.sourceTable = from_table_reference + request = bigquery.BigqueryJobsInsertRequest( projectId=project_id, job=bigquery.Job( configuration=bigquery.JobConfiguration( - copy=bigquery.JobConfigurationTableCopy( - destinationTable=to_table_reference, - sourceTable=from_table_reference, - createDisposition=create_disposition, - writeDisposition=write_disposition, - ), + copy=copy_config, labels=_build_job_labels(job_labels), ), jobReference=reference, diff --git a/sdks/python/apache_beam/version.py b/sdks/python/apache_beam/version.py index 7d1cd44ad012..870b93e530d6 100644 --- a/sdks/python/apache_beam/version.py +++ b/sdks/python/apache_beam/version.py @@ -17,4 +17,4 @@ """Apache Beam SDK version information and utilities.""" -__version__ = '2.75.0.dev' +__version__ = '2.76.0.dev' diff --git a/sdks/typescript/package.json b/sdks/typescript/package.json index a9468b04ac8b..d042bbaa7814 100644 --- a/sdks/typescript/package.json +++ b/sdks/typescript/package.json @@ -1,6 +1,6 @@ { "name": "apache-beam", - "version": "2.75.0-SNAPSHOT", + "version": "2.76.0-SNAPSHOT", "devDependencies": { "@google-cloud/bigquery": "^5.12.0", "@types/mocha": "^9.0.0", diff --git a/website/www/site/content/en/documentation/runners/flink.md b/website/www/site/content/en/documentation/runners/flink.md index d201d2546f8c..e924ccdb7bd6 100644 --- a/website/www/site/content/en/documentation/runners/flink.md +++ b/website/www/site/content/en/documentation/runners/flink.md @@ -358,12 +358,12 @@ To find out which version of Flink is compatible with Beam please see the table