From 4768e536b93358536af2f0c47565861c322b7d71 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Tue, 9 Sep 2025 08:52:42 -0700 Subject: [PATCH 01/16] fix: Improve error handling and logging during Dataproc session creation --- .../cloud/dataproc_spark_connect/session.py | 22 +++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 3d7d1dab..3a5948e2 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -450,18 +450,32 @@ def create_session_pbar(): create_session_pbar_thread.join() DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - raise DataprocSparkConnectException( + + error_msg = ( f"Error while creating Dataproc Session: {e.message}" ) + + # Only log in environments that don't auto-display exceptions + if not environment.is_colab_enterprise(): + logger.error(error_msg) + + raise DataprocSparkConnectException(error_msg) except Exception as e: stop_create_session_pbar_event.set() if create_session_pbar_thread.is_alive(): create_session_pbar_thread.join() DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - raise RuntimeError( - f"Error while creating Dataproc Session" - ) from e + + error_msg = ( + f"Error while creating Dataproc Session: {str(e)}" + ) + + # Only log in environments that don't auto-display exceptions + if not environment.is_colab_enterprise(): + logger.error(error_msg) + + raise RuntimeError(error_msg) from e finally: stop_create_session_pbar_event.set() From 14f9cf5a09a30c707fde8ae2583bd1bb451247bb Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Tue, 9 Sep 2025 09:10:33 -0700 Subject: [PATCH 02/16] update test text --- tests/unit/test_session.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/unit/test_session.py b/tests/unit/test_session.py index cd9e3514..39f8d419 100644 --- a/tests/unit/test_session.py +++ b/tests/unit/test_session.py @@ -668,7 +668,8 @@ def test_create_spark_session_with_create_session_failed( Session() ).getOrCreate() self.assertEqual( - "Error while creating Dataproc Session", e.exception.args[0] + "Error while creating Dataproc Session: Testing create session failure", + e.exception.args[0], ) @mock.patch("google.auth.default") From 9f577544130f65413ea4ad672478760eee8378b2 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Wed, 10 Sep 2025 11:08:46 -0700 Subject: [PATCH 03/16] update trace string to be array --- .../dataproc_spark_connect/exceptions.py | 2 +- .../cloud/dataproc_spark_connect/session.py | 26 +++++++------------ 2 files changed, 11 insertions(+), 17 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 53358ecf..3e5c8e90 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -24,4 +24,4 @@ def __init__(self, message): super().__init__(message) def _render_traceback_(self): - return self.message + return [self.message] diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 3a5948e2..7b8be820 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -451,15 +451,13 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - error_msg = ( - f"Error while creating Dataproc Session: {e.message}" + logger.debug(f"DEBUG: Caught InvalidArgument/PermissionDenied: {type(e)}, message: {getattr(e, 'message', 'No message attr')}") + + # Try both .message and str(e) to be safe + error_msg = getattr(e, 'message', str(e)) + raise DataprocSparkConnectException( + f"Error while creating Dataproc Session: {error_msg}" ) - - # Only log in environments that don't auto-display exceptions - if not environment.is_colab_enterprise(): - logger.error(error_msg) - - raise DataprocSparkConnectException(error_msg) except Exception as e: stop_create_session_pbar_event.set() if create_session_pbar_thread.is_alive(): @@ -467,15 +465,11 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - error_msg = ( + logger.debug(f"DEBUG: Caught other exception: {type(e)}, message: {str(e)}") + + raise RuntimeError( f"Error while creating Dataproc Session: {str(e)}" - ) - - # Only log in environments that don't auto-display exceptions - if not environment.is_colab_enterprise(): - logger.error(error_msg) - - raise RuntimeError(error_msg) from e + ) from e finally: stop_create_session_pbar_event.set() From c1c475cb83a2c23c63f31fcd0aeb4ebdb8cd0e20 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Wed, 10 Sep 2025 14:10:11 -0700 Subject: [PATCH 04/16] pyink --- google/cloud/dataproc_spark_connect/session.py | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 7b8be820..c6da09fd 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -451,12 +451,8 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - logger.debug(f"DEBUG: Caught InvalidArgument/PermissionDenied: {type(e)}, message: {getattr(e, 'message', 'No message attr')}") - - # Try both .message and str(e) to be safe - error_msg = getattr(e, 'message', str(e)) raise DataprocSparkConnectException( - f"Error while creating Dataproc Session: {error_msg}" + f"Error while creating Dataproc Session: {e.message}" ) except Exception as e: stop_create_session_pbar_event.set() @@ -465,8 +461,6 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - logger.debug(f"DEBUG: Caught other exception: {type(e)}, message: {str(e)}") - raise RuntimeError( f"Error while creating Dataproc Session: {str(e)}" ) from e From 9fab570f59add7f36f2de0857031800e6873f4ee Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 10:48:00 -0700 Subject: [PATCH 05/16] removing customized error display --- google/cloud/dataproc_spark_connect/exceptions.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 3e5c8e90..81d82583 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -23,5 +23,3 @@ def __init__(self, message): self.message = message super().__init__(message) - def _render_traceback_(self): - return [self.message] From b2cb6908eda8dc4ae4e51135dd10a8dd3d4b3af6 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 10:50:07 -0700 Subject: [PATCH 06/16] test --- google/cloud/dataproc_spark_connect/session.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index c6da09fd..5c4360cd 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -451,9 +451,9 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - raise DataprocSparkConnectException( - f"Error while creating Dataproc Session: {e.message}" - ) + error_msg = f"Error while creating Dataproc Session: {e.message}" + print(f"ABOUT TO RAISE: {error_msg}") # Debug + raise DataprocSparkConnectException(error_msg) except Exception as e: stop_create_session_pbar_event.set() if create_session_pbar_thread.is_alive(): From ee6d9602def7bbc7b78e2a8531b5993a98aea620 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 10:53:50 -0700 Subject: [PATCH 07/16] test --- google/cloud/dataproc_spark_connect/session.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 5c4360cd..ba0a103d 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -461,9 +461,9 @@ def create_session_pbar(): DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - raise RuntimeError( - f"Error while creating Dataproc Session: {str(e)}" - ) from e + error_msg = f"Error while creating Dataproc Session: {str(e)}" + print(f"SECOND HANDLER - ABOUT TO RAISE RuntimeError: {error_msg}") # Debug + raise RuntimeError(error_msg) from e finally: stop_create_session_pbar_event.set() From ce51f8c086718fa6bb4582364bc243b2759f991e Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 10:56:39 -0700 Subject: [PATCH 08/16] test --- google/cloud/dataproc_spark_connect/session.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index ba0a103d..2475177b 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -393,6 +393,7 @@ def create_session_pbar(): os.environ["SPARK_CONNECT_MODE_ENABLED"] = "1" try: + print("DEBUG: Starting session creation try block") # Debug if ( os.getenv( "DATAPROC_SPARK_CONNECT_SESSION_TERMINATE_AT_EXIT", @@ -414,6 +415,7 @@ def create_session_pbar(): self._display_session_link_on_creation(session_id) self._display_view_session_details_button(session_id) create_session_pbar_thread.start() + print("DEBUG: About to call operation.result()") # Debug session_response: Session = operation.result( polling=retry.Retry( predicate=POLLING_PREDICATE, From fd6b09054e2f26a60dd59da5aea9e49a315000cd Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 10:59:31 -0700 Subject: [PATCH 09/16] test --- google/cloud/dataproc_spark_connect/session.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 2475177b..0a00df6e 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -320,6 +320,7 @@ def __create_spark_connect_session_from_s8s( return session def __create(self) -> "DataprocSparkSession": + print("DEBUG: Entering __create() method") # Debug with self._lock: if self._options.get("spark.remote", False): @@ -329,8 +330,10 @@ def __create(self) -> "DataprocSparkSession": from google.cloud.dataproc_v1 import SessionControllerClient + print("DEBUG: About to call _get_dataproc_config()") # Debug dataproc_config: Session = self._get_dataproc_config() + print("DEBUG: About to call _check_runtime_compatibility()") # Debug # Check runtime version compatibility before creating session self._check_runtime_compatibility(dataproc_config) From b270d13322a71e378250524f67a0632fd3118fe1 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 11:01:28 -0700 Subject: [PATCH 10/16] test --- google/cloud/dataproc_spark_connect/session.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 0a00df6e..08a52183 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -320,7 +320,7 @@ def __create_spark_connect_session_from_s8s( return session def __create(self) -> "DataprocSparkSession": - print("DEBUG: Entering __create() method") # Debug + logger.error("DEBUG: Entering __create() method") # Debug with self._lock: if self._options.get("spark.remote", False): @@ -330,10 +330,10 @@ def __create(self) -> "DataprocSparkSession": from google.cloud.dataproc_v1 import SessionControllerClient - print("DEBUG: About to call _get_dataproc_config()") # Debug + logger.error("DEBUG: About to call _get_dataproc_config()") # Debug dataproc_config: Session = self._get_dataproc_config() - print("DEBUG: About to call _check_runtime_compatibility()") # Debug + logger.error("DEBUG: About to call _check_runtime_compatibility()") # Debug # Check runtime version compatibility before creating session self._check_runtime_compatibility(dataproc_config) @@ -396,7 +396,7 @@ def create_session_pbar(): os.environ["SPARK_CONNECT_MODE_ENABLED"] = "1" try: - print("DEBUG: Starting session creation try block") # Debug + logger.error("DEBUG: Starting session creation try block") # Debug if ( os.getenv( "DATAPROC_SPARK_CONNECT_SESSION_TERMINATE_AT_EXIT", @@ -418,7 +418,7 @@ def create_session_pbar(): self._display_session_link_on_creation(session_id) self._display_view_session_details_button(session_id) create_session_pbar_thread.start() - print("DEBUG: About to call operation.result()") # Debug + logger.error("DEBUG: About to call operation.result()") # Debug session_response: Session = operation.result( polling=retry.Retry( predicate=POLLING_PREDICATE, @@ -457,7 +457,7 @@ def create_session_pbar(): DataprocSparkSession._active_session_uses_custom_id = False error_msg = f"Error while creating Dataproc Session: {e.message}" - print(f"ABOUT TO RAISE: {error_msg}") # Debug + logger.error(f"ABOUT TO RAISE: {error_msg}") # Debug raise DataprocSparkConnectException(error_msg) except Exception as e: stop_create_session_pbar_event.set() @@ -467,7 +467,7 @@ def create_session_pbar(): DataprocSparkSession._active_session_uses_custom_id = False error_msg = f"Error while creating Dataproc Session: {str(e)}" - print(f"SECOND HANDLER - ABOUT TO RAISE RuntimeError: {error_msg}") # Debug + logger.error(f"SECOND HANDLER - ABOUT TO RAISE RuntimeError: {error_msg}") # Debug raise RuntimeError(error_msg) from e finally: stop_create_session_pbar_event.set() From 8e52f70eadaf46d949e87ba58417fd28ff6fdc36 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Thu, 11 Sep 2025 11:02:35 -0700 Subject: [PATCH 11/16] test --- google/cloud/dataproc_spark_connect/session.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 08a52183..04430605 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -558,14 +558,22 @@ def _get_exiting_active_session( return None def getOrCreate(self) -> "DataprocSparkSession": + logger.error("DEBUG: Entering getOrCreate()") # Debug with DataprocSparkSession._lock: + logger.error("DEBUG: Inside lock in getOrCreate()") # Debug # Handle custom session ID by setting it early and letting existing logic handle it if self._custom_session_id: + logger.error("DEBUG: Handling custom session ID") # Debug self._handle_custom_session_id() + logger.error("DEBUG: About to call _get_exiting_active_session()") # Debug session = self._get_exiting_active_session() + logger.error(f"DEBUG: _get_exiting_active_session() returned: {session is not None}") # Debug if session is None: + logger.error("DEBUG: About to call __create()") # Debug session = self.__create() + logger.error("DEBUG: __create() completed successfully") # Debug + logger.error("DEBUG: About to return session from getOrCreate()") # Debug return session def _handle_custom_session_id(self): From 683c83a830cd945a5602691abf4ace3b315557f0 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Mon, 15 Sep 2025 12:10:17 -0700 Subject: [PATCH 12/16] fixed --- .../dataproc_spark_connect/exceptions.py | 53 +++++++++++++++++++ .../cloud/dataproc_spark_connect/session.py | 34 +++++------- tests/unit/test_session.py | 3 +- 3 files changed, 67 insertions(+), 23 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 81d82583..3702cd1a 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -19,7 +19,60 @@ class DataprocSparkConnectException(Exception): doesn't provide any additional information. """ + _ipython_handler_patched = False + def __init__(self, message): self.message = message super().__init__(message) + if not DataprocSparkConnectException._ipython_handler_patched: + self._setup_ipython_exception_handler() + + def _render_traceback_(self): + return [self.message] + + def _setup_ipython_exception_handler(self): + """Setup custom exception handler for IPython environments to ensure minimal traceback display.""" + try: + from IPython import get_ipython + import sys + + ipython = get_ipython() + if ipython is not None: + # Store original method if not already stored + if not hasattr(ipython, "_original_showtraceback"): + ipython._original_showtraceback = ipython.showtraceback + + def custom_showtraceback( + shell, + exc_tuple=None, + filename=None, + tb_offset=None, + exception_only=False, + running_compiled_code=False, + ): + # Get the current exception info + _, value, _ = ( + sys.exc_info() if exc_tuple is None else exc_tuple + ) + + # If it's our custom exception, show only the message + if isinstance(value, DataprocSparkConnectException): + print(f"Error: {value.message}", file=sys.stderr) + else: + # Use original behavior for other exceptions + shell._original_showtraceback( + exc_tuple, + filename, + tb_offset, + exception_only, + running_compiled_code, + ) + + # Override the method + ipython.showtraceback = custom_showtraceback + # Mark as patched to avoid redundant setup + DataprocSparkConnectException._ipython_handler_patched = True + except ImportError: + # Not in IPython environment, no action needed + pass diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 04430605..168ee8b2 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -320,7 +320,6 @@ def __create_spark_connect_session_from_s8s( return session def __create(self) -> "DataprocSparkSession": - logger.error("DEBUG: Entering __create() method") # Debug with self._lock: if self._options.get("spark.remote", False): @@ -330,10 +329,8 @@ def __create(self) -> "DataprocSparkSession": from google.cloud.dataproc_v1 import SessionControllerClient - logger.error("DEBUG: About to call _get_dataproc_config()") # Debug dataproc_config: Session = self._get_dataproc_config() - logger.error("DEBUG: About to call _check_runtime_compatibility()") # Debug # Check runtime version compatibility before creating session self._check_runtime_compatibility(dataproc_config) @@ -396,7 +393,6 @@ def create_session_pbar(): os.environ["SPARK_CONNECT_MODE_ENABLED"] = "1" try: - logger.error("DEBUG: Starting session creation try block") # Debug if ( os.getenv( "DATAPROC_SPARK_CONNECT_SESSION_TERMINATE_AT_EXIT", @@ -418,7 +414,6 @@ def create_session_pbar(): self._display_session_link_on_creation(session_id) self._display_view_session_details_button(session_id) create_session_pbar_thread.start() - logger.error("DEBUG: About to call operation.result()") # Debug session_response: Session = operation.result( polling=retry.Retry( predicate=POLLING_PREDICATE, @@ -455,20 +450,18 @@ def create_session_pbar(): create_session_pbar_thread.join() DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - - error_msg = f"Error while creating Dataproc Session: {e.message}" - logger.error(f"ABOUT TO RAISE: {error_msg}") # Debug - raise DataprocSparkConnectException(error_msg) + raise DataprocSparkConnectException( + f"Error while creating Dataproc Session: {e.message}" + ) except Exception as e: stop_create_session_pbar_event.set() if create_session_pbar_thread.is_alive(): create_session_pbar_thread.join() DataprocSparkSession._active_s8s_session_id = None DataprocSparkSession._active_session_uses_custom_id = False - - error_msg = f"Error while creating Dataproc Session: {str(e)}" - logger.error(f"SECOND HANDLER - ABOUT TO RAISE RuntimeError: {error_msg}") # Debug - raise RuntimeError(error_msg) from e + raise RuntimeError( + f"Error while creating Dataproc Session" + ) from e finally: stop_create_session_pbar_event.set() @@ -558,22 +551,21 @@ def _get_exiting_active_session( return None def getOrCreate(self) -> "DataprocSparkSession": - logger.error("DEBUG: Entering getOrCreate()") # Debug with DataprocSparkSession._lock: - logger.error("DEBUG: Inside lock in getOrCreate()") # Debug # Handle custom session ID by setting it early and letting existing logic handle it if self._custom_session_id: - logger.error("DEBUG: Handling custom session ID") # Debug self._handle_custom_session_id() - logger.error("DEBUG: About to call _get_exiting_active_session()") # Debug session = self._get_exiting_active_session() - logger.error(f"DEBUG: _get_exiting_active_session() returned: {session is not None}") # Debug if session is None: - logger.error("DEBUG: About to call __create()") # Debug session = self.__create() - logger.error("DEBUG: __create() completed successfully") # Debug - logger.error("DEBUG: About to return session from getOrCreate()") # Debug + + # Register this session as the instantiated SparkSession for compatibility + # with tools and libraries that expect SparkSession._instantiatedSession + from pyspark.sql import SparkSession as PySparkSQLSession + + PySparkSQLSession._instantiatedSession = session + return session def _handle_custom_session_id(self): diff --git a/tests/unit/test_session.py b/tests/unit/test_session.py index 39f8d419..cd9e3514 100644 --- a/tests/unit/test_session.py +++ b/tests/unit/test_session.py @@ -668,8 +668,7 @@ def test_create_spark_session_with_create_session_failed( Session() ).getOrCreate() self.assertEqual( - "Error while creating Dataproc Session: Testing create session failure", - e.exception.args[0], + "Error while creating Dataproc Session", e.exception.args[0] ) @mock.patch("google.auth.default") From 3880b12cb23b79cb95d6a9067254f3c0e66a0335 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Mon, 15 Sep 2025 12:11:59 -0700 Subject: [PATCH 13/16] clean up code --- google/cloud/dataproc_spark_connect/session.py | 6 ------ 1 file changed, 6 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index 168ee8b2..ab4b41b1 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -560,12 +560,6 @@ def getOrCreate(self) -> "DataprocSparkSession": if session is None: session = self.__create() - # Register this session as the instantiated SparkSession for compatibility - # with tools and libraries that expect SparkSession._instantiatedSession - from pyspark.sql import SparkSession as PySparkSQLSession - - PySparkSQLSession._instantiatedSession = session - return session def _handle_custom_session_id(self): From 84777ba3644adaa85ecdd591625783cae6c92082 Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Tue, 23 Sep 2025 18:03:43 -0700 Subject: [PATCH 14/16] keep it updated with 138 --- .../dataproc_spark_connect/exceptions.py | 104 +++++++++--------- 1 file changed, 53 insertions(+), 51 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 3702cd1a..4e44259b 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -12,6 +12,59 @@ # See the License for the specific language governing permissions and # limitations under the License. +import sys + + +def _setup_ipython_exception_handler(): + """Setup custom exception handler for IPython environments to ensure minimal traceback display.""" + try: + from IPython import get_ipython + except ImportError: + return + + ipython = get_ipython() + if ipython is None: + return + + # Store original method if not already stored + if hasattr(ipython, "_dataproc_spark_connect_original_showtraceback"): + return # Already patched + + ipython._dataproc_spark_connect_original_showtraceback = ( + ipython.showtraceback + ) + + def custom_showtraceback( + shell, + exc_tuple=None, + filename=None, + tb_offset=None, + exception_only=False, + running_compiled_code=False, + ): + # Get the current exception info + _, value, _ = sys.exc_info() if exc_tuple is None else exc_tuple + + # If it's our custom exception, show only the message + if isinstance(value, DataprocSparkConnectException): + print(f"Error: {value.message}", file=sys.stderr) + else: + # Use original behavior for other exceptions + shell._dataproc_spark_connect_original_showtraceback( + exc_tuple, + filename, + tb_offset, + exception_only, + running_compiled_code, + ) + + # Override the method + ipython.showtraceback = custom_showtraceback + + +# Setup the handler once at module import time +_setup_ipython_exception_handler() + class DataprocSparkConnectException(Exception): """A custom exception class to only print the error messages. @@ -19,60 +72,9 @@ class DataprocSparkConnectException(Exception): doesn't provide any additional information. """ - _ipython_handler_patched = False - def __init__(self, message): self.message = message super().__init__(message) - if not DataprocSparkConnectException._ipython_handler_patched: - self._setup_ipython_exception_handler() def _render_traceback_(self): return [self.message] - - def _setup_ipython_exception_handler(self): - """Setup custom exception handler for IPython environments to ensure minimal traceback display.""" - try: - from IPython import get_ipython - import sys - - ipython = get_ipython() - if ipython is not None: - # Store original method if not already stored - if not hasattr(ipython, "_original_showtraceback"): - ipython._original_showtraceback = ipython.showtraceback - - def custom_showtraceback( - shell, - exc_tuple=None, - filename=None, - tb_offset=None, - exception_only=False, - running_compiled_code=False, - ): - # Get the current exception info - _, value, _ = ( - sys.exc_info() if exc_tuple is None else exc_tuple - ) - - # If it's our custom exception, show only the message - if isinstance(value, DataprocSparkConnectException): - print(f"Error: {value.message}", file=sys.stderr) - else: - # Use original behavior for other exceptions - shell._original_showtraceback( - exc_tuple, - filename, - tb_offset, - exception_only, - running_compiled_code, - ) - - # Override the method - ipython.showtraceback = custom_showtraceback - # Mark as patched to avoid redundant setup - DataprocSparkConnectException._ipython_handler_patched = True - - except ImportError: - # Not in IPython environment, no action needed - pass From c909091870ebf3511fb50f9d793d28167275548c Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Wed, 24 Sep 2025 08:27:14 -0700 Subject: [PATCH 15/16] keep consistent with 138 --- google/cloud/dataproc_spark_connect/exceptions.py | 8 +++----- google/cloud/dataproc_spark_connect/session.py | 1 - 2 files changed, 3 insertions(+), 6 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 4e44259b..69257bfc 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -27,12 +27,10 @@ def _setup_ipython_exception_handler(): return # Store original method if not already stored - if hasattr(ipython, "_dataproc_spark_connect_original_showtraceback"): + if hasattr(ipython, "_original_showtraceback"): return # Already patched - ipython._dataproc_spark_connect_original_showtraceback = ( - ipython.showtraceback - ) + ipython._original_showtraceback = ipython.showtraceback def custom_showtraceback( shell, @@ -50,7 +48,7 @@ def custom_showtraceback( print(f"Error: {value.message}", file=sys.stderr) else: # Use original behavior for other exceptions - shell._dataproc_spark_connect_original_showtraceback( + shell._original_showtraceback( exc_tuple, filename, tb_offset, diff --git a/google/cloud/dataproc_spark_connect/session.py b/google/cloud/dataproc_spark_connect/session.py index ab4b41b1..3d7d1dab 100644 --- a/google/cloud/dataproc_spark_connect/session.py +++ b/google/cloud/dataproc_spark_connect/session.py @@ -559,7 +559,6 @@ def getOrCreate(self) -> "DataprocSparkSession": session = self._get_exiting_active_session() if session is None: session = self.__create() - return session def _handle_custom_session_id(self): From c5de3c18df4a1330ce46a6ca9ff92937567b7f2c Mon Sep 17 00:00:00 2001 From: fangyh20 Date: Wed, 24 Sep 2025 08:28:43 -0700 Subject: [PATCH 16/16] updated --- google/cloud/dataproc_spark_connect/exceptions.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/google/cloud/dataproc_spark_connect/exceptions.py b/google/cloud/dataproc_spark_connect/exceptions.py index 69257bfc..4e44259b 100644 --- a/google/cloud/dataproc_spark_connect/exceptions.py +++ b/google/cloud/dataproc_spark_connect/exceptions.py @@ -27,10 +27,12 @@ def _setup_ipython_exception_handler(): return # Store original method if not already stored - if hasattr(ipython, "_original_showtraceback"): + if hasattr(ipython, "_dataproc_spark_connect_original_showtraceback"): return # Already patched - ipython._original_showtraceback = ipython.showtraceback + ipython._dataproc_spark_connect_original_showtraceback = ( + ipython.showtraceback + ) def custom_showtraceback( shell, @@ -48,7 +50,7 @@ def custom_showtraceback( print(f"Error: {value.message}", file=sys.stderr) else: # Use original behavior for other exceptions - shell._original_showtraceback( + shell._dataproc_spark_connect_original_showtraceback( exc_tuple, filename, tb_offset,