From ebd2952877c3235d19a939d2e1be0209fa948db8 Mon Sep 17 00:00:00 2001 From: pdebelak Date: Wed, 22 Feb 2023 15:59:40 -0600 Subject: [PATCH 1/3] Avoid including fallback message for S3TaskHandler if streaming logs This checks the metadata for the log_pos and if it is present and greater than 1 it doesn't include the "Falling back to local log" message. --- .../amazon/aws/log/s3_task_handler.py | 6 ++++- .../amazon/aws/log/test_s3_task_handler.py | 27 +++++++++++++++++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/airflow/providers/amazon/aws/log/s3_task_handler.py b/airflow/providers/amazon/aws/log/s3_task_handler.py index 831c86417185f..c25681e9e0376 100644 --- a/airflow/providers/amazon/aws/log/s3_task_handler.py +++ b/airflow/providers/amazon/aws/log/s3_task_handler.py @@ -128,8 +128,12 @@ def _read(self, ti, try_number, metadata=None): if logs: return "".join(f"*** {x}\n" for x in messages) + "\n".join(logs), {"end_of_log": True} else: + if metadata and "log_pos" in metadata and metadata["log_pos"] > 0: + log = "" + else: + log = "*** Falling back to local log\n" local_log, metadata = super()._read(ti, try_number, metadata) - return "*** Falling back to local log\n" + local_log, metadata + return log + local_log, metadata def s3_log_exists(self, remote_log_location: str) -> bool: """ diff --git a/tests/providers/amazon/aws/log/test_s3_task_handler.py b/tests/providers/amazon/aws/log/test_s3_task_handler.py index 4fbe453274a9b..aeca09d36d6a7 100644 --- a/tests/providers/amazon/aws/log/test_s3_task_handler.py +++ b/tests/providers/amazon/aws/log/test_s3_task_handler.py @@ -146,6 +146,33 @@ def test_read_when_s3_log_missing(self): assert actual == expected assert {"end_of_log": True, "log_pos": 0} == metadata[0] + def test_read_when_s3_log_missing_and_log_pos_missing_pre_26(self): + ti = copy.copy(self.ti) + ti.state = TaskInstanceState.SUCCESS + # mock that super class has no _read_remote_logs method + with mock.patch("airflow.providers.amazon.aws.log.s3_task_handler.hasattr", return_value=False): + log, metadata = self.s3_task_handler.read(ti) + assert 1 == len(log) + assert log[0][0][-1].startswith("*** Falling back to local log") + + def test_read_when_s3_log_missing_and_log_pos_zero_pre_26(self): + ti = copy.copy(self.ti) + ti.state = TaskInstanceState.SUCCESS + # mock that super class has no _read_remote_logs method + with mock.patch("airflow.providers.amazon.aws.log.s3_task_handler.hasattr", return_value=False): + log, metadata = self.s3_task_handler.read(ti, metadata={"log_pos": 0}) + assert 1 == len(log) + assert log[0][0][-1].startswith("*** Falling back to local log") + + def test_read_when_s3_log_missing_and_log_pos_over_zero_pre_26(self): + ti = copy.copy(self.ti) + ti.state = TaskInstanceState.SUCCESS + # mock that super class has no _read_remote_logs method + with mock.patch("airflow.providers.amazon.aws.log.s3_task_handler.hasattr", return_value=False): + log, metadata = self.s3_task_handler.read(ti, metadata={"log_pos": 1}) + assert 1 == len(log) + assert not log[0][0][-1].startswith("*** Falling back to local log") + def test_s3_read_when_log_missing(self): handler = self.s3_task_handler url = "s3://bucket/foo" From 510f4b583b32165db8a073bbb22701dc7aec507d Mon Sep 17 00:00:00 2001 From: Peter Debelak Date: Wed, 22 Feb 2023 16:58:26 -0600 Subject: [PATCH 2/3] Apply code clean-up suggestions to s3 task handler fallback handling Co-authored-by: Daniel Standish <15932138+dstandish@users.noreply.github.com> Co-authored-by: Philippe Gagnon <12717218+pgagnon@users.noreply.github.com> --- airflow/providers/amazon/aws/log/s3_task_handler.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/airflow/providers/amazon/aws/log/s3_task_handler.py b/airflow/providers/amazon/aws/log/s3_task_handler.py index c25681e9e0376..bd8ff303bc198 100644 --- a/airflow/providers/amazon/aws/log/s3_task_handler.py +++ b/airflow/providers/amazon/aws/log/s3_task_handler.py @@ -128,10 +128,10 @@ def _read(self, ti, try_number, metadata=None): if logs: return "".join(f"*** {x}\n" for x in messages) + "\n".join(logs), {"end_of_log": True} else: - if metadata and "log_pos" in metadata and metadata["log_pos"] > 0: - log = "" + if metadata and metadata.get("log_pos", 0) > 0: + log_prefix = "" else: - log = "*** Falling back to local log\n" + log_prefix = "*** Falling back to local log\n" local_log, metadata = super()._read(ti, try_number, metadata) return log + local_log, metadata From 92aac29dceb6a18147c6ebc5fa6d633b25c1513d Mon Sep 17 00:00:00 2001 From: Peter Debelak Date: Wed, 22 Feb 2023 16:59:22 -0600 Subject: [PATCH 3/3] Use fstring rather than string concatenation Co-authored-by: Philippe Gagnon <12717218+pgagnon@users.noreply.github.com> --- airflow/providers/amazon/aws/log/s3_task_handler.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/airflow/providers/amazon/aws/log/s3_task_handler.py b/airflow/providers/amazon/aws/log/s3_task_handler.py index bd8ff303bc198..098f17a28a95c 100644 --- a/airflow/providers/amazon/aws/log/s3_task_handler.py +++ b/airflow/providers/amazon/aws/log/s3_task_handler.py @@ -133,7 +133,7 @@ def _read(self, ti, try_number, metadata=None): else: log_prefix = "*** Falling back to local log\n" local_log, metadata = super()._read(ti, try_number, metadata) - return log + local_log, metadata + return f"{log_prefix}{local_log}", metadata def s3_log_exists(self, remote_log_location: str) -> bool: """