Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 25 additions & 19 deletions tests/integ/modin/test_session_cte_optimization.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@
#

import logging
import pandas as native_pd
import threading
from contextlib import contextmanager
from typing import Any, Generator
Expand All @@ -16,6 +15,10 @@

from tests.integ.utils.sql_counter import sql_count_checker

_MULTITHREAD_SESSION_WARNING = (
"You might have more than one threads sharing the Session object trying to update"
)


@contextmanager
def session_parameter_override(
Expand All @@ -31,29 +34,32 @@ def session_parameter_override(


@sql_count_checker(query_count=0)
def test_cte_optimization_for_snowpark_pandas(db_parameters, caplog):
def test_cte_optimization_for_snowpark_pandas(caplog):
session = _get_active_session()

caplog.set_level(logging.WARNING)

# Creating a dataframe ensures that another thread is running which triggers the warning
session.create_dataframe(
native_pd.DataFrame([[1, 11], [2, 12], [2, 13]], columns=["A", "B"])
)
assert len(threading.enumerate()) > 1

# Use context manager to temporarily disable CTE optimization
with session_parameter_override(session, "cte_optimization_enabled", False):
assert session.cte_optimization_enabled is False
# The multithreading warning is emitted only while another thread is alive.
stop = threading.Event()
extra_thread = threading.Thread(target=stop.wait, name="cte-opt-warning-probe")
extra_thread.start()
try:
assert len(threading.enumerate()) > 1

caplog.clear()
session.cte_optimization_enabled = session.cte_optimization_enabled
assert _MULTITHREAD_SESSION_WARNING in caplog.text

# Verify that cte optimization will be enabled once snowpark pandas is used
assert pd.session == session
assert pd.session.cte_optimization_enabled is True
with session_parameter_override(session, "cte_optimization_enabled", False):
assert session.cte_optimization_enabled is False

# Verify that the warning is filtered out
assert (
"You might have more than one threads sharing the Session object trying to update"
not in caplog.text
)
caplog.clear()

assert pd.session == session
assert pd.session.cte_optimization_enabled is True

assert _MULTITHREAD_SESSION_WARNING not in caplog.text
finally:
stop.set()
extra_thread.join(timeout=5)
assert not extra_thread.is_alive()
Loading