From 8e8e2d39dbd2db7dc4da48d5d60b7333f41f90a7 Mon Sep 17 00:00:00 2001 From: Filip Pawlowski Date: Thu, 6 Aug 2026 14:50:39 +0000 Subject: [PATCH 1/2] SNOW-2912540: inline _create_temp_stage and _create_temp_file_format into analyzer_utils Removes the dependency on these private snowflake-connector-python functions (previously imported from connector.pandas_tools). Ports the implementation from universal-driver PR #518 directly into Snowpark so the connector no longer needs to expose staging helpers as part of its public surface. Co-Authored-By: Claude Sonnet 4.6 --- .../_internal/analyzer/analyzer_utils.py | 127 +++++++++++++++++- 1 file changed, 120 insertions(+), 7 deletions(-) diff --git a/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py b/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py index cf8f701570..6be0e27380 100644 --- a/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py +++ b/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py @@ -7,17 +7,13 @@ import math import os import re +import secrets import tempfile -from typing import Any, Dict, List, Optional, Tuple, Union, Literal, Sequence +from typing import Any, Callable, Dict, List, Optional, Tuple, Union, Literal, Sequence from snowflake.connector import ProgrammingError from snowflake.connector.cursor import SnowflakeCursor from snowflake.connector.options import pyarrow -from snowflake.connector.pandas_tools import ( - _create_temp_stage, - _create_temp_file_format, - build_location_helper, -) from snowflake.snowpark._internal.analyzer.binary_plan_node import ( AsOf, Except, @@ -2212,6 +2208,123 @@ def cte_statement(queries: List[str], table_names: List[str]) -> str: return f"{WITH}{result}" +# --------------------------------------------------------------------------- +# Pandas/Arrow staging helpers +# (ported from snowflake.connector.pandas_tools so Snowpark owns the impl) +# --------------------------------------------------------------------------- + +_PANDAS_COMPRESSION_MAP: Dict[str, str] = {"gzip": "auto", "snappy": "snappy", "none": "none"} + + +def _pandas_quote_identifier(name: str) -> str: + return '"' + name.replace('"', '""') + '"' + + +def _qualify_name( + database: Optional[str], + schema: Optional[str], + name: str, + quote_identifiers: bool, +) -> str: + parts = [p for p in (database, schema, name) if p is not None] + if quote_identifiers: + parts = [_pandas_quote_identifier(p) for p in parts] + return ".".join(parts) + + +build_location_helper = _qualify_name + + +def _pandas_generate_temp_name(prefix: str) -> str: + return f"__WRITE_PANDAS_{prefix}_{secrets.token_hex(8)}" + + +def _pandas_create_temp_object( + cursor: SnowflakeCursor, + sql_builder: Callable[[str], Tuple[str, tuple]], + qualified_name: str, + bare_name: str, +) -> str: + sql, params = sql_builder(qualified_name) + try: + cursor.execute(sql, params=params, _force_qmark_paramstyle=True) + return qualified_name + except ProgrammingError: + sql, params = sql_builder(bare_name) + cursor.execute(sql, params=params, _force_qmark_paramstyle=True) + return bare_name + + +def _stage_sql( + name: str, + compression: str, + binary_as_text_false: bool, + use_scoped: bool = False, +) -> Tuple[str, tuple]: + mapped = _PANDAS_COMPRESSION_MAP[compression] + temp_type = "SCOPED TEMPORARY" if use_scoped else "TEMPORARY" + fmt_opts = [f"TYPE=PARQUET COMPRESSION={mapped}"] + if binary_as_text_false: + fmt_opts.append("BINARY_AS_TEXT=FALSE") + return (f"CREATE {temp_type} STAGE IDENTIFIER(?) FILE_FORMAT=({' '.join(fmt_opts)})", (name,)) + + +def _file_format_sql( + name: str, + compression: str, + use_logical_type_suffix: str = "", + use_scoped: bool = False, +) -> Tuple[str, tuple]: + mapped = _PANDAS_COMPRESSION_MAP[compression] + temp_type = "SCOPED TEMPORARY" if use_scoped else "TEMPORARY" + return ( + f"CREATE {temp_type} FILE FORMAT IDENTIFIER(?) TYPE=PARQUET COMPRESSION={mapped}{use_logical_type_suffix}", + (name,), + ) + + +def _create_temp_stage( + cursor: SnowflakeCursor, + database: Optional[str], + schema: Optional[str], + quote_identifiers: bool, + compression: str, + auto_create_table: bool, + overwrite: bool, + use_scoped_temp_object: bool = False, +) -> str: + name = _pandas_generate_temp_name("STAGE") + qualified = _qualify_name(database, schema, name, quote_identifiers) + return _pandas_create_temp_object( + cursor, + lambda n: _stage_sql(n, compression, auto_create_table or overwrite, use_scoped_temp_object), + qualified, + name, + ) + + +def _create_temp_file_format( + cursor: SnowflakeCursor, + database: Optional[str], + schema: Optional[str], + quote_identifiers: bool, + compression: str, + sql_use_logical_type: str, + use_scoped_temp_object: bool = False, +) -> str: + name = _pandas_generate_temp_name("FILE_FORMAT") + qualified = _qualify_name(database, schema, name, quote_identifiers) + return _pandas_create_temp_object( + cursor, + lambda n: _file_format_sql(n, compression, sql_use_logical_type, use_scoped_temp_object), + qualified, + name, + ) + + +# --------------------------------------------------------------------------- + + def write_arrow( cursor: SnowflakeCursor, table: "pyarrow.Table", @@ -2364,7 +2477,7 @@ def drop_object(name: str, object_type: str) -> None: database, schema, quote_identifiers, - compression_map[compression], + compression, sql_use_logical_type, use_scoped_temp_object, ) From cb14275c18e8df8a484861b3388a314e9b0a19ac Mon Sep 17 00:00:00 2001 From: Filip Pawlowski Date: Mon, 17 Aug 2026 18:58:09 +0000 Subject: [PATCH 2/2] SNOW-2912540: fix black formatting and pyright overload error in inlined pandas staging helpers black never ran on the ported code (4 blocks needed line-wrapping). The two cursor.execute(..., _force_qmark_paramstyle=True) calls also fail pyright: SnowflakeCursor.execute's two @overload stubs in snowflake-connector-python never declare _force_qmark_paramstyle, only the real implementation signature does. That gap was invisible while this code lived in the connector package; porting it into analyzer_utils.py exposes it to Snowpark's own pyright run. --- .../_internal/analyzer/analyzer_utils.py | 25 ++++++++++++++----- 1 file changed, 19 insertions(+), 6 deletions(-) diff --git a/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py b/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py index 6be0e27380..09511b201d 100644 --- a/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py +++ b/src/snowflake/snowpark/_internal/analyzer/analyzer_utils.py @@ -2213,7 +2213,11 @@ def cte_statement(queries: List[str], table_names: List[str]) -> str: # (ported from snowflake.connector.pandas_tools so Snowpark owns the impl) # --------------------------------------------------------------------------- -_PANDAS_COMPRESSION_MAP: Dict[str, str] = {"gzip": "auto", "snappy": "snappy", "none": "none"} +_PANDAS_COMPRESSION_MAP: Dict[str, str] = { + "gzip": "auto", + "snappy": "snappy", + "none": "none", +} def _pandas_quote_identifier(name: str) -> str: @@ -2247,11 +2251,13 @@ def _pandas_create_temp_object( ) -> str: sql, params = sql_builder(qualified_name) try: - cursor.execute(sql, params=params, _force_qmark_paramstyle=True) + # _force_qmark_paramstyle is only on execute()'s real signature, not its + # @overload stubs, so pyright can't match any overload for this call. + cursor.execute(sql, params=params, _force_qmark_paramstyle=True) # type: ignore[call-overload] return qualified_name except ProgrammingError: sql, params = sql_builder(bare_name) - cursor.execute(sql, params=params, _force_qmark_paramstyle=True) + cursor.execute(sql, params=params, _force_qmark_paramstyle=True) # type: ignore[call-overload] return bare_name @@ -2266,7 +2272,10 @@ def _stage_sql( fmt_opts = [f"TYPE=PARQUET COMPRESSION={mapped}"] if binary_as_text_false: fmt_opts.append("BINARY_AS_TEXT=FALSE") - return (f"CREATE {temp_type} STAGE IDENTIFIER(?) FILE_FORMAT=({' '.join(fmt_opts)})", (name,)) + return ( + f"CREATE {temp_type} STAGE IDENTIFIER(?) FILE_FORMAT=({' '.join(fmt_opts)})", + (name,), + ) def _file_format_sql( @@ -2297,7 +2306,9 @@ def _create_temp_stage( qualified = _qualify_name(database, schema, name, quote_identifiers) return _pandas_create_temp_object( cursor, - lambda n: _stage_sql(n, compression, auto_create_table or overwrite, use_scoped_temp_object), + lambda n: _stage_sql( + n, compression, auto_create_table or overwrite, use_scoped_temp_object + ), qualified, name, ) @@ -2316,7 +2327,9 @@ def _create_temp_file_format( qualified = _qualify_name(database, schema, name, quote_identifiers) return _pandas_create_temp_object( cursor, - lambda n: _file_format_sql(n, compression, sql_use_logical_type, use_scoped_temp_object), + lambda n: _file_format_sql( + n, compression, sql_use_logical_type, use_scoped_temp_object + ), qualified, name, )