From 950840036086180bf5ce95b0fd02185f544e999c Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Sat, 12 Feb 2022 12:12:30 +0100 Subject: [PATCH 1/9] handle logging Java and jaydebeapi exceptions when task fails --- airflow/models/taskinstance.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/airflow/models/taskinstance.py b/airflow/models/taskinstance.py index 34f5a6e2cee82..c4c3cd7eb34c3 100644 --- a/airflow/models/taskinstance.py +++ b/airflow/models/taskinstance.py @@ -47,6 +47,8 @@ import dill import jinja2 import pendulum +import jpype +from jaydebeapi import DatabaseError, InterfaceError from jinja2 import TemplateAssertionError, UndefinedError from sqlalchemy import ( Column, @@ -1721,6 +1723,9 @@ def handle_failure( test_mode = self.test_mode if error: + if jpype.isJVMStarted(): + if isinstance(error, (jpype.java.sql.SQLException, DatabaseError, InterfaceError)): + self.log.error("%s", error) if isinstance(error, BaseException): self.log.error("Task failed with exception", exc_info=error) else: From 7c0c38f5bb9ae99b19d111971b2b05250c4ba942 Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Sat, 12 Feb 2022 16:40:32 +0100 Subject: [PATCH 2/9] Catch AttributeError when logging BaseException error --- airflow/models/taskinstance.py | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/airflow/models/taskinstance.py b/airflow/models/taskinstance.py index c4c3cd7eb34c3..69caeab6fbc23 100644 --- a/airflow/models/taskinstance.py +++ b/airflow/models/taskinstance.py @@ -47,8 +47,6 @@ import dill import jinja2 import pendulum -import jpype -from jaydebeapi import DatabaseError, InterfaceError from jinja2 import TemplateAssertionError, UndefinedError from sqlalchemy import ( Column, @@ -1723,11 +1721,11 @@ def handle_failure( test_mode = self.test_mode if error: - if jpype.isJVMStarted(): - if isinstance(error, (jpype.java.sql.SQLException, DatabaseError, InterfaceError)): - self.log.error("%s", error) if isinstance(error, BaseException): - self.log.error("Task failed with exception", exc_info=error) + try: + self.log.error("Task failed with exception", exc_info=error) + except AttributeError: + self.log.error("%s", error) else: self.log.error("%s", error) # external monitoring process provides pickle file so _run_raw_task From a26e0dbef851729c51539541c1e48975ac16331f Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Thu, 17 Feb 2022 19:40:27 +0100 Subject: [PATCH 3/9] catch AttributeError when exception has read only properties --- airflow/models/taskinstance.py | 5 +---- airflow/utils/log/secrets_masker.py | 5 ++++- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/airflow/models/taskinstance.py b/airflow/models/taskinstance.py index 69caeab6fbc23..34f5a6e2cee82 100644 --- a/airflow/models/taskinstance.py +++ b/airflow/models/taskinstance.py @@ -1722,10 +1722,7 @@ def handle_failure( if error: if isinstance(error, BaseException): - try: - self.log.error("Task failed with exception", exc_info=error) - except AttributeError: - self.log.error("%s", error) + self.log.error("Task failed with exception", exc_info=error) else: self.log.error("%s", error) # external monitoring process provides pickle file so _run_raw_task diff --git a/airflow/utils/log/secrets_masker.py b/airflow/utils/log/secrets_masker.py index 538057b13ae1d..665fd9113ec08 100644 --- a/airflow/utils/log/secrets_masker.py +++ b/airflow/utils/log/secrets_masker.py @@ -147,7 +147,10 @@ def _record_attrs_to_ignore(self) -> Iterable[str]: return frozenset(record.__dict__).difference({'msg', 'args'}) def _redact_exception_with_context(self, exception): - exception.args = (self.redact(v) for v in exception.args) + try: + exception.args = (self.redact(v) for v in exception.args) + except AttributeError: + pass if exception.__context__: self._redact_exception_with_context(exception.__context__) if exception.__cause__ and exception.__cause__ is not exception.__context__: From b098466f7cd04173da563c8ed185ffe0a3ab2633 Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Sun, 20 Feb 2022 15:28:12 +0100 Subject: [PATCH 4/9] check Oracle connection using dual table --- airflow/providers/oracle/hooks/oracle.py | 13 +++++++++++++ tests/providers/oracle/hooks/test_oracle.py | 4 ++++ 2 files changed, 17 insertions(+) diff --git a/airflow/providers/oracle/hooks/oracle.py b/airflow/providers/oracle/hooks/oracle.py index 6843d8f5c4469..7cbdde3ad7ee4 100644 --- a/airflow/providers/oracle/hooks/oracle.py +++ b/airflow/providers/oracle/hooks/oracle.py @@ -339,3 +339,16 @@ def handler(cursor): ) return result + + def test_connection(self): + """Tests the connection by executing a select 1 from dual query""" + status, message = False, '' + try: + if self.get_first("select 1 from dual"): + status = True + message = 'Connection successfully tested' + except Exception as e: + status = False + message = str(e) + + return status, message diff --git a/tests/providers/oracle/hooks/test_oracle.py b/tests/providers/oracle/hooks/test_oracle.py index c837c0d7db99b..0da63523f3a37 100644 --- a/tests/providers/oracle/hooks/test_oracle.py +++ b/tests/providers/oracle/hooks/test_oracle.py @@ -347,3 +347,7 @@ def bindvar(value): expected = [1, 0, 0.0, False, ''] assert self.cur.execute.mock_calls == [mock.call('BEGIN proc(:1,:2,:3,:4,:5); END;', expected)] assert result == expected + + def test_test_connection_use_dual_table(self): + self.db_hook.test_connection() + self.cur.execute.assert_called_once_with("select 1 from dual") From d584ef72601d345fdb1214e8005c0d43d561faaf Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Sun, 20 Feb 2022 15:38:31 +0100 Subject: [PATCH 5/9] undo commit from different issue --- airflow/providers/oracle/hooks/oracle.py | 13 ------------- tests/providers/oracle/hooks/test_oracle.py | 4 ---- 2 files changed, 17 deletions(-) diff --git a/airflow/providers/oracle/hooks/oracle.py b/airflow/providers/oracle/hooks/oracle.py index 7cbdde3ad7ee4..6843d8f5c4469 100644 --- a/airflow/providers/oracle/hooks/oracle.py +++ b/airflow/providers/oracle/hooks/oracle.py @@ -339,16 +339,3 @@ def handler(cursor): ) return result - - def test_connection(self): - """Tests the connection by executing a select 1 from dual query""" - status, message = False, '' - try: - if self.get_first("select 1 from dual"): - status = True - message = 'Connection successfully tested' - except Exception as e: - status = False - message = str(e) - - return status, message diff --git a/tests/providers/oracle/hooks/test_oracle.py b/tests/providers/oracle/hooks/test_oracle.py index 0da63523f3a37..c837c0d7db99b 100644 --- a/tests/providers/oracle/hooks/test_oracle.py +++ b/tests/providers/oracle/hooks/test_oracle.py @@ -347,7 +347,3 @@ def bindvar(value): expected = [1, 0, 0.0, False, ''] assert self.cur.execute.mock_calls == [mock.call('BEGIN proc(:1,:2,:3,:4,:5); END;', expected)] assert result == expected - - def test_test_connection_use_dual_table(self): - self.db_hook.test_connection() - self.cur.execute.assert_called_once_with("select 1 from dual") From e34afd03d24e62f30026f7fc1900776e96b54267 Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Mon, 21 Feb 2022 06:42:35 +0100 Subject: [PATCH 6/9] added comment explaining catching error --- airflow/utils/log/secrets_masker.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/airflow/utils/log/secrets_masker.py b/airflow/utils/log/secrets_masker.py index 665fd9113ec08..27f4bb1920a7a 100644 --- a/airflow/utils/log/secrets_masker.py +++ b/airflow/utils/log/secrets_masker.py @@ -147,6 +147,8 @@ def _record_attrs_to_ignore(self) -> Iterable[str]: return frozenset(record.__dict__).difference({'msg', 'args'}) def _redact_exception_with_context(self, exception): + # In case, when exception has "read only" properties, we need + # to catch an AttributeError to log properly try: exception.args = (self.redact(v) for v in exception.args) except AttributeError: From 2ac3a379d8c85d9b7c000f54715cd8e049a25125 Mon Sep 17 00:00:00 2001 From: hubert-pietron <94397721+hubert-pietron@users.noreply.github.com> Date: Mon, 21 Feb 2022 07:59:24 +0100 Subject: [PATCH 7/9] Update airflow/utils/log/secrets_masker.py Co-authored-by: Tzu-ping Chung --- airflow/utils/log/secrets_masker.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/airflow/utils/log/secrets_masker.py b/airflow/utils/log/secrets_masker.py index 27f4bb1920a7a..c3f20aa5b166f 100644 --- a/airflow/utils/log/secrets_masker.py +++ b/airflow/utils/log/secrets_masker.py @@ -147,8 +147,8 @@ def _record_attrs_to_ignore(self) -> Iterable[str]: return frozenset(record.__dict__).difference({'msg', 'args'}) def _redact_exception_with_context(self, exception): - # In case, when exception has "read only" properties, we need - # to catch an AttributeError to log properly + # Exception class may not be modifiable (e.g. declared by an + # extension module such as JDBC). try: exception.args = (self.redact(v) for v in exception.args) except AttributeError: From 3f9073c46940ef1f25a9f46b447d9cf84435c3ed Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Tue, 3 May 2022 17:37:19 +0200 Subject: [PATCH 8/9] prevent closing connection using chunksize parameter in get_pandas_df --- airflow/hooks/dbapi.py | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/airflow/hooks/dbapi.py b/airflow/hooks/dbapi.py index 1f898706a5e91..49aa50046c928 100644 --- a/airflow/hooks/dbapi.py +++ b/airflow/hooks/dbapi.py @@ -120,6 +120,9 @@ def get_pandas_df(self, sql, parameters=None, **kwargs): :param parameters: The parameters to render the SQL query with. :param kwargs: (optional) passed into pandas.io.sql.read_sql method """ + if 'chunksize' in kwargs: + return self.get_pandas_df_by_chunks(sql, parameters=parameters, **kwargs) + try: from pandas.io import sql as psql except ImportError: @@ -128,6 +131,15 @@ def get_pandas_df(self, sql, parameters=None, **kwargs): with closing(self.get_conn()) as conn: return psql.read_sql(sql, con=conn, params=parameters, **kwargs) + def get_pandas_df_by_chunks(self, sql, parameters=None, **kwargs): + try: + from pandas.io import sql as psql + except ImportError: + raise Exception("pandas library not installed, run: pip install 'apache-airflow[pandas]'.") + + with closing(self.get_conn()) as conn: + yield from psql.read_sql(sql, con=conn, params=parameters, **kwargs) + def get_records(self, sql, parameters=None): """ Executes the sql and returns a set of records. From f524ad5403fca9944a62798ad7021b26580ae9ce Mon Sep 17 00:00:00 2001 From: hubert-pietron Date: Sat, 7 May 2022 10:56:22 +0200 Subject: [PATCH 9/9] make get_pandas_df_by_chunks separate function --- airflow/hooks/dbapi.py | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/airflow/hooks/dbapi.py b/airflow/hooks/dbapi.py index 49aa50046c928..97250100a9576 100644 --- a/airflow/hooks/dbapi.py +++ b/airflow/hooks/dbapi.py @@ -120,9 +120,6 @@ def get_pandas_df(self, sql, parameters=None, **kwargs): :param parameters: The parameters to render the SQL query with. :param kwargs: (optional) passed into pandas.io.sql.read_sql method """ - if 'chunksize' in kwargs: - return self.get_pandas_df_by_chunks(sql, parameters=parameters, **kwargs) - try: from pandas.io import sql as psql except ImportError: @@ -131,14 +128,23 @@ def get_pandas_df(self, sql, parameters=None, **kwargs): with closing(self.get_conn()) as conn: return psql.read_sql(sql, con=conn, params=parameters, **kwargs) - def get_pandas_df_by_chunks(self, sql, parameters=None, **kwargs): + def get_pandas_df_by_chunks(self, sql, parameters=None, *, chunksize, **kwargs): + """ + Executes the sql and returns a generator + + :param sql: the sql statement to be executed (str) or a list of + sql statements to execute + :param parameters: The parameters to render the SQL query with + :param chunksize: number of rows to include in each chunk + :param kwargs: (optional) passed into pandas.io.sql.read_sql method + """ try: from pandas.io import sql as psql except ImportError: raise Exception("pandas library not installed, run: pip install 'apache-airflow[pandas]'.") with closing(self.get_conn()) as conn: - yield from psql.read_sql(sql, con=conn, params=parameters, **kwargs) + yield from psql.read_sql(sql, con=conn, params=parameters, chunksize=chunksize, **kwargs) def get_records(self, sql, parameters=None): """