From ce1d10795bee07a9c145135c6d72e829276449ee Mon Sep 17 00:00:00 2001 From: David Caron Date: Wed, 20 Jan 2021 12:10:02 -0500 Subject: [PATCH 1/2] Fix returning last log line in docker operator Iterable values in `cli.attach(..., stream=True)` are not necessarily single lines. Lines can be merged together or split because of buffering. When this split happens just before the last newline character, the returned value of `execute` is an empty byte string. --- airflow/providers/docker/operators/docker.py | 5 ++++- tests/providers/docker/operators/test_docker.py | 4 ++-- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/airflow/providers/docker/operators/docker.py b/airflow/providers/docker/operators/docker.py index 0440d21b30c04..10a4f627e3074 100644 --- a/airflow/providers/docker/operators/docker.py +++ b/airflow/providers/docker/operators/docker.py @@ -269,7 +269,10 @@ def _run_image(self) -> Optional[str]: # duplicated conditional logic because of expensive operation ret = None if self.do_xcom_push: - ret = self.cli.logs(container=self.container['Id']) if self.xcom_all else line.encode('utf-8') + if self.xcom_all: + ret = self.cli.logs(container=self.container['Id']) + else: + ret = self.cli.logs(container=self.container['Id'], tail=1).strip() if self.auto_remove: self.cli.remove_container(self.container['Id']) diff --git a/tests/providers/docker/operators/test_docker.py b/tests/providers/docker/operators/test_docker.py index 0a2f8383e825f..e5ec762ac7bad 100644 --- a/tests/providers/docker/operators/test_docker.py +++ b/tests/providers/docker/operators/test_docker.py @@ -41,8 +41,8 @@ def setUp(self): self.client_mock = mock.Mock(spec=APIClient) self.client_mock.create_container.return_value = {'Id': 'some_id'} self.client_mock.images.return_value = [] - self.client_mock.attach.return_value = ['container log'] - self.client_mock.logs.return_value = ['container log'] + self.client_mock.attach.return_value = [b'container log'] + self.client_mock.logs.return_value = b'container log' # logs(..., stream=False) returns bytes self.client_mock.pull.return_value = {"status": "pull log"} self.client_mock.wait.return_value = {"StatusCode": 0} self.client_mock.create_host_config.return_value = mock.Mock() From 36367bc393ee20ccde0f11dd6008a8b511c295e3 Mon Sep 17 00:00:00 2001 From: David Caron Date: Wed, 20 Jan 2021 12:36:32 -0500 Subject: [PATCH 2/2] Fix DockerOperator.execute return type --- airflow/providers/docker/operators/docker.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/airflow/providers/docker/operators/docker.py b/airflow/providers/docker/operators/docker.py index 10a4f627e3074..6bdeddf44ad03 100644 --- a/airflow/providers/docker/operators/docker.py +++ b/airflow/providers/docker/operators/docker.py @@ -219,7 +219,7 @@ def get_hook(self) -> DockerHook: tls=self.__get_tls_config(), ) - def _run_image(self) -> Optional[str]: + def _run_image(self) -> Optional[bytes]: """Run a Docker container with the provided image""" self.log.info('Starting docker container from image %s', self.image) @@ -279,7 +279,7 @@ def _run_image(self) -> Optional[str]: return ret - def execute(self, context) -> Optional[str]: + def execute(self, context) -> Optional[bytes]: self.cli = self._get_cli() if not self.cli: raise Exception("The 'cli' should be initialized before!")