This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new a048842a72a Fix google provider system tests dag version increase
issues (#72194)
a048842a72a is described below
commit a048842a72a7ef1490da3cb984611f3c78474491
Author: olegkachur-e <[email protected]>
AuthorDate: Mon Oct 5 16:39:56 2026 +0000
Fix google provider system tests dag version increase issues (#72194)
* Fix google provider system tests dag version increase issues
- Fix DAG version incrementing on every parse cycle by replacing dynamic
values.
* Keep the Composer and Cloud Tasks system test Dags stable across parses
CloudComposerExternalTaskSensor does not template execution_range, so
Jinja strings reached poke() unrendered; a one-day timedelta gives the
same window. The Cloud Tasks schedule_time still came from datetime.now()
at parse time, so it now comes from a task at run time.
Generated-by: Claude Opus 5
---------
Co-authored-by: Oleg Kachur <[email protected]>
Co-authored-by: Jarek Potiuk <[email protected]>
---
.../google/cloud/bigquery/example_bigquery_dts.py | 3 +-
.../cloud/bigquery/example_bigquery_sensors.py | 3 -
.../cloud/bigquery/example_bigquery_tables.py | 3 +-
.../cloud/bigquery/example_bigquery_value_check.py | 2 +-
.../cloud/composer/example_cloud_composer.py | 4 +-
.../google/cloud/dataplex/example_dataplex_dp.py | 2 +-
.../example_cloud_storage_transfer_service_aws.py | 64 ++++++++++++++--------
.../example_cloud_storage_transfer_service_gcp.py | 55 ++++++++++++-------
.../system/google/cloud/tasks/example_tasks.py | 15 +++--
.../marketing_platform/example_campaign_manager.py | 6 +-
10 files changed, 92 insertions(+), 65 deletions(-)
diff --git
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_dts.py
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_dts.py
index c4f698d1281..400deaeac42 100644
---
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_dts.py
+++
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_dts.py
@@ -22,7 +22,6 @@ Example Airflow DAG that creates and deletes Bigquery data
transfer configuratio
from __future__ import annotations
import os
-import time
from datetime import datetime
from pathlib import Path
from typing import cast
@@ -136,7 +135,7 @@ with DAG(
task_id="gcp_bigquery_start_transfer",
project_id=PROJECT_ID,
transfer_config_id=transfer_config_id,
- requested_run_time={"seconds": int(time.time() + 60)},
+ requested_run_time={"seconds": "{{ (macros.datetime.now().timestamp()
| int) + 60 }}"},
)
# [END howto_bigquery_start_transfer]
diff --git
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_sensors.py
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_sensors.py
index 59d5b210715..1e5ae09967e 100644
---
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_sensors.py
+++
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_sensors.py
@@ -48,10 +48,7 @@ DAG_ID = "bigquery_sensors"
DATASET_NAME = f"dataset_{DAG_ID}_{ENV_ID}".replace("-", "_")
TABLE_NAME = f"partitioned_table_{DAG_ID}_{ENV_ID}".replace("-", "_")
-
-INSERT_DATE = datetime.now().strftime("%Y-%m-%d")
PARTITION_NAME = "{{ ds_nodash }}"
-
INSERT_ROWS_QUERY = f"INSERT {DATASET_NAME}.{TABLE_NAME} VALUES (42, '{{{{ ds
}}}}')"
SCHEMA = [
diff --git
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
index b1fca85e409..4a2814d0590 100644
---
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
+++
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
@@ -22,7 +22,6 @@ Example Airflow DAG for Google BigQuery service testing
tables.
from __future__ import annotations
import os
-import time
from datetime import datetime
from pathlib import Path
@@ -178,7 +177,7 @@ with DAG(
dataset_id=DATASET_NAME,
table_resource={
"tableReference": {"tableId": "test_table_id"},
- "expirationTime": (int(time.time()) + 300) * 1000,
+ "expirationTime": "{{ (macros.datetime.now().timestamp() | int +
300) * 1000 }}",
},
)
# [END howto_operator_bigquery_upsert_table]
diff --git
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_value_check.py
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_value_check.py
index dd5fbf8f252..705690236b5 100644
---
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_value_check.py
+++
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_value_check.py
@@ -55,7 +55,7 @@ SCHEMA = [
DAG_ID = "bq_value_check_location"
DATASET = f"ds_{DAG_ID}_{ENV_ID}"
TABLE = "ds_table"
-INSERT_DATE = datetime.now().strftime("%Y-%m-%d")
+INSERT_DATE = "{{ ds }}"
INSERT_ROWS_QUERY = (
f"INSERT {DATASET}.{TABLE} VALUES "
f"(42, 'monty python', '{INSERT_DATE}'), "
diff --git
a/providers/google/tests/system/google/cloud/composer/example_cloud_composer.py
b/providers/google/tests/system/google/cloud/composer/example_cloud_composer.py
index 5077d5fb311..5b73858068e 100644
---
a/providers/google/tests/system/google/cloud/composer/example_cloud_composer.py
+++
b/providers/google/tests/system/google/cloud/composer/example_cloud_composer.py
@@ -310,7 +310,7 @@ with DAG(
composer_external_dag_id="airflow_monitoring",
composer_external_task_id="echo",
allowed_states=["success"],
- execution_range=[datetime.now() - timedelta(1), datetime.now()],
+ execution_range=timedelta(days=1),
)
# [END howto_sensor_external_task]
@@ -323,7 +323,7 @@ with DAG(
composer_external_dag_id="airflow_monitoring",
composer_external_task_id="echo",
allowed_states=["success"],
- execution_range=[datetime.now() - timedelta(1), datetime.now()],
+ execution_range=timedelta(days=1),
deferrable=True,
)
# [END howto_sensor_external_task_deferrable_mode]
diff --git
a/providers/google/tests/system/google/cloud/dataplex/example_dataplex_dp.py
b/providers/google/tests/system/google/cloud/dataplex/example_dataplex_dp.py
index 5d315b67ffc..7548348ecad 100644
--- a/providers/google/tests/system/google/cloud/dataplex/example_dataplex_dp.py
+++ b/providers/google/tests/system/google/cloud/dataplex/example_dataplex_dp.py
@@ -77,7 +77,7 @@ SCHEMA = [
{"name": "dt", "type": "STRING", "mode": "NULLABLE"},
]
-INSERT_DATE = datetime.now().strftime("%Y-%m-%d")
+INSERT_DATE = "{{ ds }}"
INSERT_ROWS_QUERY = f"INSERT {DATASET}.{TABLE_1} VALUES (1, 'test test2',
'{INSERT_DATE}');"
LOCATION = "us"
diff --git
a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py
b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py
index aa94eec8134..cadde91390d 100644
---
a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py
+++
b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_aws.py
@@ -66,9 +66,11 @@ from
airflow.providers.google.cloud.sensors.cloud_storage_transfer_service impor
)
try:
- from airflow.sdk import TriggerRule
+ from airflow.sdk import TriggerRule, task
+
except ImportError:
# Compatibility for Airflow < 3.1
+ from airflow.decorators import task # type: ignore[no-redef,attr-defined]
from airflow.utils.trigger_rule import TriggerRule # type:
ignore[no-redef,attr-defined]
from system.google import DEFAULT_GCP_SYSTEM_TEST_PROJECT_ID
@@ -88,27 +90,30 @@ GCP_DESCRIPTION = "description"
GCP_TRANSFER_JOB_NAME =
f"transferJobs/sampleJob-{DAG_ID}-{ENV_ID}".replace("-", "_")
GCP_TRANSFER_JOB_2_NAME =
f"transferJobs/sampleJob2-{DAG_ID}-{ENV_ID}".replace("-", "_")
+
# [START howto_operator_gcp_transfer_create_job_body_aws]
-aws_to_gcs_transfer_body = {
- DESCRIPTION: GCP_DESCRIPTION,
- STATUS: GcpTransferJobsStatus.ENABLED,
- PROJECT_ID: GCP_PROJECT_ID,
- JOB_NAME: GCP_TRANSFER_JOB_NAME,
- SCHEDULE: {
- SCHEDULE_START_DATE: datetime(2015, 1, 1).date(),
- SCHEDULE_END_DATE: datetime(2030, 1, 1).date(),
- START_TIME_OF_DAY: (datetime.now(tz=timezone.utc) +
timedelta(minutes=1)).time(),
- },
- TRANSFER_SPEC: {
- AWS_S3_DATA_SOURCE: {BUCKET_NAME: BUCKET_SOURCE_AWS},
- GCS_DATA_SINK: {BUCKET_NAME: BUCKET_TARGET_GCS},
- TRANSFER_OPTIONS: {ALREADY_EXISTING_IN_SINK: True},
- },
-}
-# [END howto_operator_gcp_transfer_create_job_body_aws]
+def generate_base_transfer_body() -> dict:
+ """Helper function to generate a standard payload template dynamically at
execution time."""
+ now = datetime.now(tz=timezone.utc) + timedelta(minutes=1)
+
+ return {
+ DESCRIPTION: GCP_DESCRIPTION,
+ STATUS: GcpTransferJobsStatus.ENABLED,
+ PROJECT_ID: GCP_PROJECT_ID,
+ SCHEDULE: {
+ SCHEDULE_START_DATE: {"year": 2015, "month": 1, "day": 1},
+ SCHEDULE_END_DATE: {"year": 2030, "month": 1, "day": 1},
+ START_TIME_OF_DAY: {"hours": now.hour, "minutes": now.minute,
"seconds": now.second, "nanos": 0},
+ },
+ TRANSFER_SPEC: {
+ AWS_S3_DATA_SOURCE: {BUCKET_NAME: BUCKET_SOURCE_AWS},
+ GCS_DATA_SINK: {BUCKET_NAME: BUCKET_TARGET_GCS},
+ TRANSFER_OPTIONS: {ALREADY_EXISTING_IN_SINK: True},
+ },
+ }
-aws_to_gcs_transfer_body_2 = deepcopy(aws_to_gcs_transfer_body)
-aws_to_gcs_transfer_body_2[JOB_NAME] = GCP_TRANSFER_JOB_2_NAME
+
+# [END howto_operator_gcp_transfer_create_job_body_aws]
# [START howto_operator_gcp_transfer_update_job_body_aws]
update_body = {
@@ -132,6 +137,18 @@ with DAG(
catchup=False,
tags=["example", "aws", "gcs", "transfer"],
) as dag:
+
+ @task
+ def prepare_transfer_payloads():
+ base_transfer_body = generate_base_transfer_body()
+ transfer_body_1 = deepcopy(base_transfer_body)
+ transfer_body_1[JOB_NAME] = GCP_TRANSFER_JOB_NAME
+ transfer_body_2 = deepcopy(base_transfer_body)
+ transfer_body_2[JOB_NAME] = GCP_TRANSFER_JOB_2_NAME
+ return {"body_1": transfer_body_1, "body_2": transfer_body_2}
+
+ transfer_payloads = prepare_transfer_payloads()
+
create_bucket_s3 = S3CreateBucketOperator(
task_id="create_bucket_s3", bucket_name=BUCKET_SOURCE_AWS,
region_name="us-east-1"
)
@@ -153,7 +170,8 @@ with DAG(
# [START howto_operator_gcp_transfer_create_job]
create_transfer_job_s3_to_gcs = CloudDataTransferServiceCreateJobOperator(
- task_id="create_transfer_job_s3_to_gcs", body=aws_to_gcs_transfer_body
+ task_id="create_transfer_job_s3_to_gcs",
+ body=transfer_payloads["body_1"],
)
# [END howto_operator_gcp_transfer_create_job]
@@ -214,7 +232,8 @@ with DAG(
# [END howto_operator_gcp_transfer_update_job]
create_second_transfer_job_from_aws =
CloudDataTransferServiceCreateJobOperator(
- task_id="create_transfer_job_s3_to_gcs_2",
body=aws_to_gcs_transfer_body_2
+ task_id="create_transfer_job_s3_to_gcs_2",
+ body=transfer_payloads["body_2"],
)
wait_for_operation_to_start_2 = CloudDataTransferServiceJobStatusSensor(
@@ -265,6 +284,7 @@ with DAG(
(
# TEST SETUP
[create_bucket_s3 >> upload_file_to_s3, create_bucket_gcs]
+ >> transfer_payloads
# TEST BODY
>> create_transfer_job_s3_to_gcs
>> wait_for_operation_to_start
diff --git
a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py
b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py
index 908b7ae63b7..930c756613c 100644
---
a/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py
+++
b/providers/google/tests/system/google/cloud/storage_transfer/example_cloud_storage_transfer_service_gcp.py
@@ -1,4 +1,3 @@
-#
# 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
@@ -63,9 +62,10 @@ from
airflow.providers.google.cloud.sensors.cloud_storage_transfer_service impor
from airflow.providers.google.cloud.transfers.local_to_gcs import
LocalFilesystemToGCSOperator
try:
- from airflow.sdk import TriggerRule
+ from airflow.sdk import TriggerRule, task
except ImportError:
# Compatibility for Airflow < 3.1
+ from airflow.decorators import task # type: ignore[no-redef,attr-defined]
from airflow.utils.trigger_rule import TriggerRule # type:
ignore[no-redef,attr-defined]
ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
@@ -81,23 +81,6 @@ FILE_URI = f"gs://{BUCKET_NAME_SRC}/{FILE_NAME}"
CURRENT_FOLDER = Path(__file__).parent
FILE_LOCAL_PATH = str(Path(CURRENT_FOLDER) / "resources" / FILE_NAME)
-# [START howto_operator_gcp_transfer_create_job_body_gcp]
-gcs_to_gcs_transfer_body = {
- DESCRIPTION: "description",
- STATUS: GcpTransferJobsStatus.ENABLED,
- PROJECT_ID: PROJECT_ID_TRANSFER,
- SCHEDULE: {
- SCHEDULE_START_DATE: datetime(2015, 1, 1).date(),
- SCHEDULE_END_DATE: datetime(2030, 1, 1).date(),
- START_TIME_OF_DAY: (datetime.now(tz=timezone.utc) +
timedelta(seconds=120)).time(),
- },
- TRANSFER_SPEC: {
- GCS_DATA_SOURCE: {BUCKET_NAME: BUCKET_NAME_SRC},
- GCS_DATA_SINK: {BUCKET_NAME: BUCKET_NAME_DST},
- TRANSFER_OPTIONS: {ALREADY_EXISTING_IN_SINK: True},
- },
-}
-# [END howto_operator_gcp_transfer_create_job_body_gcp]
# [START howto_operator_gcp_transfer_update_job_body]
update_body = {
@@ -114,6 +97,37 @@ with DAG(
catchup=False,
tags=["example", "transfer", "gcp"],
) as dag:
+ # [START howto_operator_gcp_transfer_create_job_body_gcp]
+ @task
+ def prepare_transfer_payload() -> dict:
+ """Generates the payload dynamically right before job creation."""
+ now_utc = datetime.now(tz=timezone.utc) + timedelta(seconds=120)
+
+ return {
+ DESCRIPTION: "description",
+ STATUS: GcpTransferJobsStatus.ENABLED,
+ PROJECT_ID: PROJECT_ID_TRANSFER,
+ SCHEDULE: {
+ SCHEDULE_START_DATE: {"year": 2015, "month": 1, "day": 1},
+ SCHEDULE_END_DATE: {"year": 2030, "month": 1, "day": 1},
+ START_TIME_OF_DAY: {
+ "hours": now_utc.hour,
+ "minutes": now_utc.minute,
+ "seconds": now_utc.second,
+ "nanos": 0,
+ },
+ },
+ TRANSFER_SPEC: {
+ GCS_DATA_SOURCE: {BUCKET_NAME: BUCKET_NAME_SRC},
+ GCS_DATA_SINK: {BUCKET_NAME: BUCKET_NAME_DST},
+ TRANSFER_OPTIONS: {ALREADY_EXISTING_IN_SINK: True},
+ },
+ }
+
+ # [END howto_operator_gcp_transfer_create_job_body_gcp]
+
+ transfer_payload = prepare_transfer_payload()
+
create_bucket_src = GCSCreateBucketOperator(
task_id="create_bucket_src",
bucket_name=BUCKET_NAME_SRC,
@@ -135,7 +149,7 @@ with DAG(
create_transfer = CloudDataTransferServiceCreateJobOperator(
task_id="create_transfer",
- body=gcs_to_gcs_transfer_body,
+ body=transfer_payload,
)
# [START howto_operator_gcp_transfer_update_job]
@@ -199,6 +213,7 @@ with DAG(
(
[create_bucket_src, create_bucket_dst]
>> upload_file
+ >> transfer_payload
>> create_transfer
>> [wait_for_transfer, wait_for_transfer_defered]
>> update_transfer
diff --git a/providers/google/tests/system/google/cloud/tasks/example_tasks.py
b/providers/google/tests/system/google/cloud/tasks/example_tasks.py
index a8dda823555..9bb0834113b 100644
--- a/providers/google/tests/system/google/cloud/tasks/example_tasks.py
+++ b/providers/google/tests/system/google/cloud/tasks/example_tasks.py
@@ -23,11 +23,10 @@ runs and deletes Tasks in the Google Cloud Tasks service in
the Google Cloud.
from __future__ import annotations
import os
-from datetime import datetime, timedelta
+from datetime import datetime, timedelta, timezone
from google.api_core.retry import Retry
from google.cloud.tasks_v2.types import Queue
-from google.protobuf import timestamp_pb2
from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS
@@ -57,9 +56,6 @@ except ImportError:
ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
DAG_ID = "cloud_tasks_tasks"
-timestamp = timestamp_pb2.Timestamp()
-timestamp.FromDatetime(datetime.now() + timedelta(hours=12))
-
LOCATION = "us-central1"
# queue cannot use recent names even if queue was removed
QUEUE_ID = f"queue-{ENV_ID}-{DAG_ID.replace('_', '-')}"
@@ -68,9 +64,7 @@ TASK = {
"http_request": {
"http_method": "POST",
"url": "http://www.example.com/example",
- "body": b"",
},
- "schedule_time": timestamp,
}
with DAG(
@@ -111,11 +105,16 @@ with DAG(
)
delete_queue.trigger_rule = TriggerRule.ALL_DONE
+ @task(task_id="build_task")
+ def build_task_with_schedule_time():
+ # Computed when the task runs: a value computed when the file is
parsed changes the Dag on every parse.
+ return {**TASK, "schedule_time": datetime.now(tz=timezone.utc) +
timedelta(hours=12)}
+
# [START create_task]
create_task = CloudTasksTaskCreateOperator(
location=LOCATION,
queue_name=QUEUE_ID + "{{
task_instance.xcom_pull(task_ids='random_string') }}",
- task=TASK,
+ task=build_task_with_schedule_time(),
task_name=TASK_NAME + "{{
task_instance.xcom_pull(task_ids='random_string') }}",
retry=Retry(maximum=10.0),
timeout=5,
diff --git
a/providers/google/tests/system/google/marketing_platform/example_campaign_manager.py
b/providers/google/tests/system/google/marketing_platform/example_campaign_manager.py
index f5531185efe..43c86fc0354 100644
---
a/providers/google/tests/system/google/marketing_platform/example_campaign_manager.py
+++
b/providers/google/tests/system/google/marketing_platform/example_campaign_manager.py
@@ -29,8 +29,6 @@ from __future__ import annotations
import json
import logging
import os
-import time
-import uuid
from datetime import datetime
from typing import Any, cast
@@ -122,7 +120,7 @@ CONVERSION = {
"ordinal": "0",
"quantity": 42,
"value": 123.4,
- "timestampMicros": int(time.time()) * 1000000,
+ "timestampMicros": "{{ (macros.datetime.now().timestamp() * 1000000) | int
}}",
"customVariables": [
{
"kind": "dfareporting#customFloodlightVariable",
@@ -231,7 +229,7 @@ with DAG(
# [END howto_campaign_manager_wait_for_operation]
# [START howto_campaign_manager_get_report_operator]
- report_name = f"reports/report_{str(uuid.uuid1())}"
+ report_name = "reports/report_{{ macros.uuid.uuid1() }}"
get_report = GoogleCampaignManagerDownloadReportOperator(
task_id="get_report",
profile_id=USER_PROFILE_ID,