diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_locked_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_locked_dag.py new file mode 100644 index 0000000000000..42420c199afc0 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_locked_dag.py @@ -0,0 +1,57 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +A locked conf rule cannot be overridden by authoring. + +Required global conf: `{"scope": {"dag_id": "openlineage_policy_conf_locked_dag"}, "locked": true, "controls": {"include_source_code": false}}` +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.bash import BashOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + +DAG_ID = "openlineage_policy_conf_locked_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_locked = BashOperator(task_id="t_locked_blocked", bash_command="echo locked_test") + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + t_locked >> check_events + +extend_global_openlineage_emission_policy(t_locked, include_source_code=True) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_source_code_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_source_code_dag.py new file mode 100644 index 0000000000000..8900f11fcf2e8 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_conf_source_code_dag.py @@ -0,0 +1,59 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +An unlocked conf rule suppresses source code; authoring re-enables it per-task. + +Required global conf: `{"scope": {"dag_id": "openlineage_policy_conf_source_code_dag"}, "controls": {"include_source_code": false}}` +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.bash import BashOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + +DAG_ID = "openlineage_policy_conf_source_code_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_conf_suppressed = BashOperator(task_id="t_conf_suppressed", bash_command="echo conf_test") + t_auth_override = extend_global_openlineage_emission_policy( + BashOperator(task_id="t_authoring_override", bash_command="echo auth_override"), + include_source_code=True, + ) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + t_conf_suppressed >> t_auth_override >> check_events + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_emit_false_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_emit_false_dag.py new file mode 100644 index 0000000000000..11f8b9d64254d --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_emit_false_dag.py @@ -0,0 +1,61 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Dag-level ``emit=False`` silences everything (task + Dag events). +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.python import PythonOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + + +def _say_hello(): + print("hello") + + +DAG_ID = "openlineage_policy_dag_emit_false_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + task_a = PythonOperator(task_id="task_a", python_callable=_say_hello) + task_b = PythonOperator(task_id="task_b", python_callable=_say_hello) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + task_a >> task_b >> check_events + +extend_global_openlineage_emission_policy(dag, emit=False) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_events_false_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_events_false_dag.py new file mode 100644 index 0000000000000..00342319c52e4 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_events_false_dag.py @@ -0,0 +1,68 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Dag-level ``emit_dag_events=False`` suppresses DAG-run events only; task events +must remain fully intact (source code, inputs/outputs, ...). The task registers +hook lineage so the emitted task event carries real inputs/outputs to assert on. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.python import PythonOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + + +def _register_hook_lineage(): + from airflow.lineage.hook import get_hook_lineage_collector + + collector = get_hook_lineage_collector() + collector.add_input_asset(context=None, uri="file://host1/in1.txt") + collector.add_input_asset(context=None, uri="file://host1/in2.txt") + collector.add_output_asset(context=None, uri="file://host2/out1.txt") + collector.add_output_asset(context=None, uri="file://host2/out2.txt") + + +DAG_ID = "openlineage_policy_dag_events_false_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + run = PythonOperator(task_id="run", python_callable=_register_hook_lineage) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + run >> check_events + +extend_global_openlineage_emission_policy(dag, emit_dag_events=False) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_override_task_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_override_task_dag.py new file mode 100644 index 0000000000000..b326e49330143 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_dag_override_task_dag.py @@ -0,0 +1,57 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +A dag-level flag (``include_source_code=False``) is overridden per-task. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.bash import BashOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + +DAG_ID = "openlineage_policy_dag_override_task_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_default = BashOperator(task_id="t_default", bash_command="echo dag_default") + t_override = BashOperator(task_id="t_override", bash_command="echo task_override") + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + t_default >> t_override >> check_events + +extend_global_openlineage_emission_policy(dag, include_source_code=False) +extend_global_openlineage_emission_policy(t_override, include_source_code=True) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_extract_metadata_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_extract_metadata_dag.py new file mode 100644 index 0000000000000..112638a587ba8 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_extract_metadata_dag.py @@ -0,0 +1,74 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Per-task ``extract_operator_metadata=False`` gives empty inputs/outputs. +""" + +from __future__ import annotations + +from datetime import datetime + +from openlineage.client.event_v2 import Dataset + +from airflow import DAG +from airflow.providers.common.compat.sdk import BaseOperator +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.openlineage.extractors.base import OperatorLineage + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + + +class _PolicyCustomOp(BaseOperator): + """Custom operator returning fixed inputs/outputs via the OL extractor API.""" + + def execute(self, context): + print("custom op ran") + + def get_openlineage_facets_on_complete(self, task_instance): + return OperatorLineage( + inputs=[Dataset("s3://ol-test-in", "in.csv")], + outputs=[Dataset("s3://ol-test-out", "out.csv")], + ) + + +DAG_ID = "openlineage_policy_extract_metadata_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_with_extract = _PolicyCustomOp(task_id="t_with_extract") + t_no_extract = extend_global_openlineage_emission_policy( + _PolicyCustomOp(task_id="t_no_extract"), + extract_operator_metadata=False, + ) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + t_with_extract >> t_no_extract >> check_events + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_full_task_info_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_full_task_info_dag.py new file mode 100644 index 0000000000000..1bf264af68b93 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_full_task_info_dag.py @@ -0,0 +1,55 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Dag-level ``include_full_task_info=True`` adds extra fields in airflow.task. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.bash import BashOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + +DAG_ID = "openlineage_policy_full_task_info_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + run = BashOperator(task_id="run", bash_command="echo full_info_test") + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + run >> check_events + +extend_global_openlineage_emission_policy(dag, include_full_task_info=True) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_hook_lineage_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_hook_lineage_dag.py new file mode 100644 index 0000000000000..721672808b85d --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_hook_lineage_dag.py @@ -0,0 +1,107 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Per-task ``hook_lineage=False`` blocks asset + SQL child events. + +Also verifies precedence: ``extract_operator_metadata=False`` short-circuits the +whole extraction pipeline, so ``hook_lineage=True`` becomes inert (no asset +inputs/outputs, no SQL child events) on those tasks. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.python import PythonOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + + +def _register_hook_lineage(): + from airflow.lineage.hook import get_hook_lineage_collector + + collector = get_hook_lineage_collector() + collector.add_input_asset(context=None, uri="file://host1/in1.txt") + collector.add_input_asset(context=None, uri="file://host1/in2.txt") + collector.add_output_asset(context=None, uri="file://host2/out1.txt") + collector.add_output_asset(context=None, uri="file://host2/out2.txt") + + +def _register_sql_hook_lineage(): + from airflow.lineage.hook import get_hook_lineage_collector + from airflow.providers.common.sql.hooks.lineage import SqlJobHookLineageExtra + + collector = get_hook_lineage_collector() + collector.add_extra( + context=None, + key=SqlJobHookLineageExtra.KEY.value, + value={SqlJobHookLineageExtra.VALUE__SQL_STATEMENT.value: "SELECT 1 AS test_col"}, + ) + + +DAG_ID = "openlineage_policy_hook_lineage_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_with_hook = PythonOperator(task_id="t_with_hook", python_callable=_register_hook_lineage) + t_no_hook = extend_global_openlineage_emission_policy( + PythonOperator(task_id="t_no_hook", python_callable=_register_hook_lineage), + hook_lineage=False, + ) + t_with_sql = PythonOperator(task_id="t_with_sql_hook", python_callable=_register_sql_hook_lineage) + t_no_sql = extend_global_openlineage_emission_policy( + PythonOperator(task_id="t_no_sql_hook", python_callable=_register_sql_hook_lineage), + hook_lineage=False, + ) + t_hook_extract_false = extend_global_openlineage_emission_policy( + PythonOperator(task_id="t_hook_extract_false", python_callable=_register_hook_lineage), + extract_operator_metadata=False, + hook_lineage=True, + ) + t_sql_extract_false = extend_global_openlineage_emission_policy( + PythonOperator(task_id="t_sql_extract_false", python_callable=_register_sql_hook_lineage), + extract_operator_metadata=False, + hook_lineage=True, + ) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + ( + t_with_hook + >> t_no_hook + >> t_with_sql + >> t_no_sql + >> t_hook_extract_false + >> t_sql_extract_false + >> check_events + ) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_source_code_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_source_code_dag.py new file mode 100644 index 0000000000000..7ccff61e7a44c --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_source_code_dag.py @@ -0,0 +1,66 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Per-task ``include_source_code=False`` removes the source code job facet. + +Also verifies precedence: ``extract_operator_metadata=False`` short-circuits the +whole extraction pipeline, so ``include_source_code=True`` becomes inert (no +source code job facet) on that task. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.bash import BashOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + +DAG_ID = "openlineage_policy_source_code_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + t_with_src = BashOperator(task_id="t_with_src", bash_command="echo test_src") + t_no_src = extend_global_openlineage_emission_policy( + BashOperator(task_id="t_no_src", bash_command="echo test_no_src"), + include_source_code=False, + ) + t_extract_false = extend_global_openlineage_emission_policy( + BashOperator(task_id="t_extract_false", bash_command="echo test_extract_false"), + extract_operator_metadata=False, + include_source_code=True, + ) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + t_with_src >> t_no_src >> t_extract_false >> check_events + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/example_openlineage_policy_task_emit_false_dag.py b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_task_emit_false_dag.py new file mode 100644 index 0000000000000..19569c0f82b7f --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/example_openlineage_policy_task_emit_false_dag.py @@ -0,0 +1,61 @@ +# 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. +""" +DAG exercising the OpenLineage emission policy authoring API. + +Per-task ``emit=False`` silences task events only. +""" + +from __future__ import annotations + +from datetime import datetime + +from airflow import DAG +from airflow.providers.openlineage.api.emission_policy import extend_global_openlineage_emission_policy +from airflow.providers.standard.operators.python import PythonOperator + +from system.openlineage.expected_events import get_expected_event_file_path +from system.openlineage.operator import OpenLineageTestOperator + + +def _say_hello(): + print("hello") + + +DAG_ID = "openlineage_policy_task_emit_false_dag" + +with DAG( + dag_id=DAG_ID, + start_date=datetime(2021, 1, 1), + schedule=None, + catchup=False, + default_args={"retries": 0}, +) as dag: + task_normal = PythonOperator(task_id="task_normal", python_callable=_say_hello) + task_silenced = PythonOperator(task_id="task_silenced", python_callable=_say_hello) + check_events = OpenLineageTestOperator( + task_id="check_events", file_path=get_expected_event_file_path(DAG_ID) + ) + task_normal >> task_silenced >> check_events + +extend_global_openlineage_emission_policy(task_silenced, emit=False) + + +from tests_common.test_utils.system_tests import get_test_run # noqa: E402 + +# Needed to run the example DAG with pytest (see: contributing-docs/testing/system_tests.rst) +test_run = get_test_run(dag) diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_locked_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_locked_dag.json new file mode 100644 index 0000000000000..092e13d103285 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_locked_dag.json @@ -0,0 +1,26 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_conf_locked_dag" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_conf_locked_dag.t_locked_blocked", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_conf_locked_dag.t_locked_blocked", + "facets": { + "sourceCode": null + } + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_source_code_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_source_code_dag.json new file mode 100644 index 0000000000000..49fbdf4a27dd6 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_conf_source_code_dag.json @@ -0,0 +1,50 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_conf_source_code_dag" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_conf_source_code_dag.t_conf_suppressed", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_conf_source_code_dag.t_conf_suppressed", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_conf_source_code_dag.t_authoring_override", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo auth_override" + } + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_conf_source_code_dag.t_authoring_override", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo auth_override" + } + } + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_emit_false_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_emit_false_dag.json new file mode 100644 index 0000000000000..55f7d4b5fbec9 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_emit_false_dag.json @@ -0,0 +1,65 @@ +[ + { + "$absent": "dag-level emit=False suppresses DAG-run events", + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_emit_false_dag" + } + }, + { + "$absent": "dag-level emit=False suppresses DAG-run events", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_emit_false_dag" + } + }, + { + "$absent": "dag-level emit=False suppresses DAG-run events", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_dag_emit_false_dag" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_a" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_a" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_a" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_b" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_b" + } + }, + { + "$absent": "dag-level emit=False propagated to all tasks", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_dag_emit_false_dag.task_b" + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_events_false_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_events_false_dag.json new file mode 100644 index 0000000000000..9f2ef0472ade2 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_events_false_dag.json @@ -0,0 +1,69 @@ +[ + { + "$absent": "dag-level emit_dag_events=False suppresses DAG-run events", + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_events_false_dag" + } + }, + { + "$absent": "dag-level emit_dag_events=False suppresses DAG-run events", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_events_false_dag" + } + }, + { + "$absent": "dag-level emit_dag_events=False suppresses DAG-run events", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_dag_events_false_dag" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_dag_events_false_dag.run", + "facets": { + "sourceCode": { + "language": "python", + "sourceCode": "{{ 'add_input_asset' in result }}" + } + } + } + }, + { + "eventType": "COMPLETE", + "inputs": [ + { + "namespace": "file://host1", + "name": "/in1.txt" + }, + { + "namespace": "file://host1", + "name": "/in2.txt" + } + ], + "outputs": [ + { + "namespace": "file://host2", + "name": "/out1.txt" + }, + { + "namespace": "file://host2", + "name": "/out2.txt" + } + ], + "job": { + "name": "openlineage_policy_dag_events_false_dag.run", + "facets": { + "sourceCode": { + "language": "python", + "sourceCode": "{{ 'add_input_asset' in result }}" + } + } + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_override_task_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_override_task_dag.json new file mode 100644 index 0000000000000..3f617cb059712 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_dag_override_task_dag.json @@ -0,0 +1,50 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_override_task_dag" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_override_task_dag.t_default", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_override_task_dag.t_default", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_dag_override_task_dag.t_override", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo task_override" + } + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_dag_override_task_dag.t_override", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo task_override" + } + } + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_extract_metadata_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_extract_metadata_dag.json new file mode 100644 index 0000000000000..dc2ce7d03fd34 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_extract_metadata_dag.json @@ -0,0 +1,50 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_extract_metadata_dag" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_extract_metadata_dag.t_with_extract" + } + }, + { + "eventType": "COMPLETE", + "inputs": [ + { + "namespace": "s3://ol-test-in", + "name": "in.csv" + } + ], + "outputs": [ + { + "namespace": "s3://ol-test-out", + "name": "out.csv" + } + ], + "job": { + "name": "openlineage_policy_extract_metadata_dag.t_with_extract" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_extract_metadata_dag.t_no_extract" + } + }, + { + "eventType": "COMPLETE", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_extract_metadata_dag.t_no_extract" + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_full_task_info_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_full_task_info_dag.json new file mode 100644 index 0000000000000..09b025759255a --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_full_task_info_dag.json @@ -0,0 +1,38 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_full_task_info_dag" + } + }, + { + "eventType": "START", + "run": { + "facets": { + "airflow": { + "task": { + "bash_command": "echo full_info_test" + } + } + } + }, + "job": { + "name": "openlineage_policy_full_task_info_dag.run" + } + }, + { + "eventType": "COMPLETE", + "run": { + "facets": { + "airflow": { + "task": { + "bash_command": "echo full_info_test" + } + } + } + }, + "job": { + "name": "openlineage_policy_full_task_info_dag.run" + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_hook_lineage_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_hook_lineage_dag.json new file mode 100644 index 0000000000000..af3d3e5ea95b0 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_hook_lineage_dag.json @@ -0,0 +1,174 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_hook" + } + }, + { + "eventType": "COMPLETE", + "inputs": [ + { + "namespace": "file://host1", + "name": "/in1.txt" + }, + { + "namespace": "file://host1", + "name": "/in2.txt" + } + ], + "outputs": [ + { + "namespace": "file://host2", + "name": "/out1.txt" + }, + { + "namespace": "file://host2", + "name": "/out2.txt" + } + ], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_hook" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_hook" + } + }, + { + "eventType": "COMPLETE", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_hook" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_sql_hook" + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_sql_hook" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_sql_hook.query.1", + "facets": { + "sql": { + "query": "SELECT 1 AS test_col" + } + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_with_sql_hook.query.1", + "facets": { + "sql": { + "query": "SELECT 1 AS test_col" + } + } + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_sql_hook" + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_sql_hook" + } + }, + { + "$absent": "hook_lineage=False suppresses SQL child query events", + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_sql_hook.query.1" + } + }, + { + "$absent": "hook_lineage=False suppresses SQL child query events", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_sql_hook.query.1" + } + }, + { + "$absent": "hook_lineage=False suppresses SQL child query events", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_no_sql_hook.query.1" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_hook_extract_false" + } + }, + { + "eventType": "COMPLETE", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_hook_extract_false" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_sql_extract_false" + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_sql_extract_false" + } + }, + { + "$absent": "extract_operator_metadata=False skips all extraction incl. SQL child events", + "eventType": "START", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_sql_extract_false.query.1" + } + }, + { + "$absent": "extract_operator_metadata=False skips all extraction incl. SQL child events", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_sql_extract_false.query.1" + } + }, + { + "$absent": "extract_operator_metadata=False skips all extraction incl. SQL child events", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_hook_lineage_dag.t_sql_extract_false.query.1" + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_source_code_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_source_code_dag.json new file mode 100644 index 0000000000000..2351ffc214298 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_source_code_dag.json @@ -0,0 +1,68 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_source_code_dag" + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_source_code_dag.t_with_src", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo test_src" + } + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_source_code_dag.t_with_src", + "facets": { + "sourceCode": { + "language": "bash", + "sourceCode": "echo test_src" + } + } + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_source_code_dag.t_no_src", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_source_code_dag.t_no_src", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "START", + "job": { + "name": "openlineage_policy_source_code_dag.t_extract_false", + "facets": { + "sourceCode": null + } + } + }, + { + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_source_code_dag.t_extract_false", + "facets": { + "sourceCode": null + } + } + } +] diff --git a/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_task_emit_false_dag.json b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_task_emit_false_dag.json new file mode 100644 index 0000000000000..cde484931d318 --- /dev/null +++ b/providers/openlineage/tests/system/openlineage/expected_events/openlineage_policy_task_emit_false_dag.json @@ -0,0 +1,57 @@ +[ + { + "eventType": "START", + "job": { + "name": "openlineage_policy_task_emit_false_dag" + } + }, + { + "eventType": "START", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_task_emit_false_dag.task_normal", + "facets": { + "sourceCode": { + "language": "python", + "sourceCode": "{{ 'hello' in result }}" + } + } + } + }, + { + "eventType": "COMPLETE", + "inputs": [], + "outputs": [], + "job": { + "name": "openlineage_policy_task_emit_false_dag.task_normal", + "facets": { + "sourceCode": { + "language": "python", + "sourceCode": "{{ 'hello' in result }}" + } + } + } + }, + { + "$absent": "task has emit=False via authoring", + "eventType": "START", + "job": { + "name": "openlineage_policy_task_emit_false_dag.task_silenced" + } + }, + { + "$absent": "task has emit=False via authoring", + "eventType": "COMPLETE", + "job": { + "name": "openlineage_policy_task_emit_false_dag.task_silenced" + } + }, + { + "$absent": "task has emit=False via authoring", + "eventType": "FAIL", + "job": { + "name": "openlineage_policy_task_emit_false_dag.task_silenced" + } + } +]