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 2874061932b Revert DMS provision checks to __init__ with is not None
polarity (#70634)
2874061932b is described below
commit 2874061932be696bddcec39f193b95289b7d6afb
Author: ahilashsasidharan <[email protected]>
AuthorDate: Fri Jul 31 21:21:42 2026 -0400
Revert DMS provision checks to __init__ with is not None polarity (#70634)
* Revert DMS provision checks to __init__ with is not None polarity
* Change DmsStartReplicationOperator to raise ValueError over
AirflowException
---
generated/known_airflow_exceptions.txt | 2 +-
.../airflow/providers/amazon/aws/operators/dms.py | 11 +++----
.../tests/unit/amazon/aws/operators/test_dms.py | 36 ++++++++++------------
3 files changed, 21 insertions(+), 28 deletions(-)
diff --git a/generated/known_airflow_exceptions.txt
b/generated/known_airflow_exceptions.txt
index b99a61c7807..e170208c367 100644
--- a/generated/known_airflow_exceptions.txt
+++ b/generated/known_airflow_exceptions.txt
@@ -69,7 +69,7 @@
providers/amazon/src/airflow/providers/amazon/aws/operators/batch.py::5
providers/amazon/src/airflow/providers/amazon/aws/operators/bedrock.py::6
providers/amazon/src/airflow/providers/amazon/aws/operators/comprehend.py::2
providers/amazon/src/airflow/providers/amazon/aws/operators/datasync.py::10
-providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py::6
+providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py::5
providers/amazon/src/airflow/providers/amazon/aws/operators/ec2.py::1
providers/amazon/src/airflow/providers/amazon/aws/operators/ecs.py::7
providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py::12
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py
b/providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py
index 1cdd9884467..ecf83773421 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/operators/dms.py
@@ -195,6 +195,8 @@ class DmsModifyTaskOperator(AwsBaseOperator[DmsHook]):
**kwargs,
):
super().__init__(aws_conn_id=aws_conn_id, **kwargs)
+ if cdc_start_time is not None and cdc_start_position is not None:
+ raise ValueError("Only one of cdc_start_time or cdc_start_position
can be provided.")
self.replication_task_arn = replication_task_arn
self.table_mappings = table_mappings
self.migration_type = migration_type
@@ -215,9 +217,6 @@ class DmsModifyTaskOperator(AwsBaseOperator[DmsHook]):
)
def execute(self, context: Context) -> dict:
- if self.cdc_start_time and self.cdc_start_position:
- raise ValueError("Only one of cdc_start_time or cdc_start_position
can be provided.")
-
tasks = self.hook.find_replication_tasks_by_arn(
replication_task_arn=self.replication_task_arn,
without_settings=True
)
@@ -919,7 +918,8 @@ class DmsStartReplicationOperator(AwsBaseOperator[DmsHook]):
aws_conn_id=aws_conn_id,
**kwargs,
)
-
+ if cdc_start_time is not None and cdc_start_pos is not None:
+ raise ValueError("Only one of cdc_start_time or cdc_start_pos
should be provided.")
self.replication_config_arn = replication_config_arn
self.replication_start_type = replication_start_type
self.cdc_start_time = cdc_start_time
@@ -931,9 +931,6 @@ class DmsStartReplicationOperator(AwsBaseOperator[DmsHook]):
self.wait_for_completion = wait_for_completion
def execute(self, context: Context):
- if self.cdc_start_time and self.cdc_start_pos:
- raise AirflowException("Only one of cdc_start_time or
cdc_start_pos should be provided.")
-
result = self.hook.describe_replications(
filters=[{"Name": "replication-config-arn", "Values":
[self.replication_config_arn]}]
)
diff --git a/providers/amazon/tests/unit/amazon/aws/operators/test_dms.py
b/providers/amazon/tests/unit/amazon/aws/operators/test_dms.py
index 2c2a299353d..bb48f7b4f66 100644
--- a/providers/amazon/tests/unit/amazon/aws/operators/test_dms.py
+++ b/providers/amazon/tests/unit/amazon/aws/operators/test_dms.py
@@ -209,16 +209,14 @@ class TestDmsModifyTaskOperator:
def _modifying_task(self):
return [{"ReplicationTaskArn": self.TASK_ARN, "Status": "modifying"}]
- def test_execute_raises_if_both_cdc_start_params_provided(self):
- op = DmsModifyTaskOperator(
- task_id="modify_task",
- replication_task_arn=self.TASK_ARN,
- cdc_start_time=datetime(2024, 1, 1),
- cdc_start_position="mysql-bin.000001:4",
- )
-
+ def test_init_raises_if_both_cdc_start_params_provided(self):
with pytest.raises(ValueError, match="Only one of"):
- op.execute(None)
+ DmsModifyTaskOperator(
+ task_id="modify_task",
+ replication_task_arn=self.TASK_ARN,
+ cdc_start_time=datetime(2024, 1, 1),
+ cdc_start_position="mysql-bin.000001:4",
+ )
@pytest.mark.parametrize("status", ["stopped", "ready", "failed"])
@mock.patch.object(DmsHook, "find_replication_tasks_by_arn")
@@ -1445,17 +1443,15 @@ class TestDmsStartReplicationOperator:
}
}
- def test_execute_raises_if_both_cdc_start_params_provided(self):
- op = DmsStartReplicationOperator(
- task_id="start_replication",
- replication_config_arn="XXXXXXXXXXXXXXX",
- replication_start_type="cdc",
- cdc_start_pos=1,
- cdc_start_time="2024-01-01 00:00:00",
- )
-
- with pytest.raises(AirflowException, match="Only one of"):
- op.execute({})
+ def test_init_raises_if_both_cdc_start_params_provided(self):
+ with pytest.raises(ValueError, match="Only one of"):
+ DmsStartReplicationOperator(
+ task_id="start_replication",
+ replication_config_arn="XXXXXXXXXXXXXXX",
+ replication_start_type="cdc",
+ cdc_start_pos=1,
+ cdc_start_time="2024-01-01 00:00:00",
+ )
@mock.patch.object(DmsHook, "describe_replications")
@mock.patch.object(DmsHook, "start_replication")