§
    šŠtj`  ã                  ó`   — d dl mZ d dlmZmZmZmZmZ er
d dlm	Z	m
Z
mZ  G d„ d¦  «        ZdS )é    )Úannotations)ÚTYPE_CHECKINGÚAnyÚIterableÚListÚOptional)Ú	DataFrameÚRowÚSparkSessionc                  óœ   — e Zd ZdZ	 	 	 	 	 	 d-d.d„Ze	 d/d0d„¦   «         Zd1d„Zd1d„Zd2d„Z	d/d3d„Z
d2d„Zd4d"„Zd5d&„Zd6d7d*„Zd/d3d+„Zd6d7d,„ZdS )8ÚSparkSQLz;SparkSQL is a utility class for interacting with Spark SQL.Né   Úspark_sessionúOptional[SparkSession]ÚcatalogúOptional[str]ÚschemaÚignore_tablesúOptional[List[str]]Úinclude_tablesÚsample_rows_in_table_infoÚintc                óX  — 	 ddl m} n# t          $ r t          d¦  «        ‚w xY w|r|n|j                             ¦   «         | _        |�| j        j                             |¦  «         |�| j        j                             |¦  «         t          |  
                    ¦   «         ¦  «        | _        |rt          |¦  «        nt          ¦   «         | _        | j        r$| j        | j        z
  }|rt          d|› d�¦  «        ‚|rt          |¦  «        nt          ¦   «         | _        | j        r$| j        | j        z
  }|rt          d|› d�¦  «        ‚|                      ¦   «         }	|	rt          |	¦  «        n| j        | _        t#          |t$          ¦  «        st'          d¦  «        ‚|| _        dS )	aÁ  Initialize a SparkSQL object.

        Args:
            spark_session: A SparkSession object.
              If not provided, one will be created.
            catalog: The catalog to use.
              If not provided, the default catalog will be used.
            schema: The schema to use.
              If not provided, the default schema will be used.
            ignore_tables: A list of tables to ignore.
              If not provided, all tables will be used.
            include_tables: A list of tables to include.
              If not provided, all tables will be used.
            sample_rows_in_table_info: The number of rows to include in the table info.
              Defaults to 3.
        r   ©r   úFpyspark is not installed. Please install it with `pip install pyspark`Nzinclude_tables ú not found in databasezignore_tables z,sample_rows_in_table_info must be an integer)Úpyspark.sqlr   ÚImportErrorÚbuilderÚgetOrCreateÚ_sparkr   ÚsetCurrentCatalogÚsetCurrentDatabaseÚsetÚ_get_all_table_namesÚ_all_tablesÚ_include_tablesÚ
ValueErrorÚ_ignore_tablesÚget_usable_table_namesÚ_usable_tablesÚ
isinstancer   Ú	TypeErrorÚ_sample_rows_in_table_info)
Úselfr   r   r   r   r   r   r   Úmissing_tablesÚusable_tabless
             úe/var/www/html/CA-Chatbot/venv/lib/python3.11/site-packages/langchain_community/utilities/spark_sql.pyÚ__init__zSparkSQL.__init__   sò  € ð2	Ø0Ð0Ð0Ð0Ð0Ð0Ð0øÝð 	ð 	ð 	ÝØXñô ð ð	øøøð +ÐRˆMˆM°Ô0D×0PÒ0PÑ0RÔ0Rð 	Œð ÐØŒKÔ×1Ò1°'Ñ:Ô:Ð:ØÐØŒKÔ×2Ò2°6Ñ:Ô:Ð:å˜t×8Ò8Ñ:Ô:Ñ;Ô;ˆÔØ6DÐO�s >Ñ2Ô2Ð2Í#É%Ì%ˆÔØÔð 	Ø!Ô1°DÔ4DÑDˆNØð Ý ØL nÐLÐLÐLñô ð ð 5BÐL�c -Ñ0Ô0Ð0ÅsÁuÄuˆÔØÔð 	Ø!Ô0°4Ô3CÑCˆNØð Ý ØK ^ÐKÐKÐKñô ð ð ×3Ò3Ñ5Ô5ˆØ4AÐW�c -Ñ0Ô0Ð0ÀtÔGWˆÔåÐ3µSÑ9Ô9ð 	LÝÐJÑKÔKÐKà*CˆÔ'Ð'Ð'ó   ‚	 ‰#Údatabase_uriÚstrÚengine_argsúOptional[dict]Úkwargsr   Úreturnc                ó¶   — 	 ddl m} n# t          $ r t          d¦  «        ‚w xY w|j                             |¦  «                             ¦   «         } | |fi |¤ŽS )zzCreating a remote Spark Session via Spark connect.
        For example: SparkSQL.from_uri("sc://localhost:15002")
        r   r   r   )r   r   r   r   Úremoter    )Úclsr5   r7   r9   r   Úsparks         r2   Úfrom_urizSparkSQL.from_uriK   sˆ   € ð	Ø0Ð0Ð0Ð0Ð0Ð0Ð0øÝð 	ð 	ð 	ÝØXñô ð ð	øøøð
 Ô$×+Ò+¨LÑ9Ô9×EÒEÑGÔGˆØˆs�5Ð#Ð#˜FÐ#Ð#Ð#r4   úIterable[str]c                óV   — | j         r| j         S t          | j        | j        z
  ¦  «        S )zGet names of tables available.)r'   Úsortedr&   r)   )r/   s    r2   r*   zSparkSQL.get_usable_table_names\   s/   € àÔð 	(ØÔ'Ð'å�dÔ&¨Ô)<Ñ<Ñ=Ô=Ð=ó    c                ó¼   — | j                              d¦  «                             d¦  «                             ¦   «         }t	          t          d„ |¦  «        ¦  «        S )NzSHOW TABLESÚ	tableNamec                ó   — | j         S ©N)rE   )Úrows    r2   ú<lambda>z/SparkSQL._get_all_table_names.<locals>.<lambda>e   s   €  C¤M€ rC   )r!   ÚsqlÚselectÚcollectÚlistÚmap)r/   Úrowss     r2   r%   zSparkSQL._get_all_table_namesc   sK   € ØŒ{�Š˜}Ñ-Ô-×4Ò4°[ÑAÔA×IÒIÑKÔKˆÝ•CÐ1Ð1°4Ñ8Ô8Ñ9Ô9Ð9rC   Útablec                óº   — | j                              d|› �¦  «                             ¦   «         d         j        }|                     d¦  «        }|d |…         dz   S )NzSHOW CREATE TABLE r   ÚUSINGú;)r!   rJ   rL   Úcreatetab_stmtÚfind)r/   rP   Ú	statementÚusing_clause_indexs       r2   Ú_get_create_table_stmtzSparkSQL._get_create_table_stmtg   s`   € àŒK�OŠOÐ8°Ð8Ð8Ñ9Ô9×AÒAÑCÔCÀAÔFÔUð 	ð 'Ÿ^š^¨GÑ4Ô4ÐØÐ,Ð,Ð,Ô-°Ñ3Ð3rC   Útable_namesc                óŠ  — |                       ¦   «         }|�9t          |¦  «                             |¦  «        }|rt          d|› d�¦  «        ‚|}g }|D ]Y}|                      |¦  «        }| j        r&|dz  }|d|                      |¦  «        › d�z  }|dz  }|                     |¦  «         ŒZd                     |¦  «        }|S )Nztable_names r   z

/*ú
z*/z

)	r*   r$   Ú
differencer(   rX   r.   Ú_get_sample_spark_rowsÚappendÚjoin)r/   rY   Úall_table_namesr0   ÚtablesÚ
table_nameÚ
table_infoÚ	final_strs           r2   Úget_table_infozSparkSQL.get_table_infoo   sò   € Ø×5Ò5Ñ7Ô7ˆØÐ"Ý  Ñ-Ô-×8Ò8¸ÑIÔIˆNØð XÝ Ð!V°Ð!VÐ!VÐ!VÑWÔWÐWØ)ˆOØˆØ)ð 	&ð 	&ˆJØ×4Ò4°ZÑ@Ô@ˆJØÔ.ð #Ø˜hÑ&�
ØÐN 4×#>Ò#>¸zÑ#JÔ#JÐNÐNÐNÑN�
Ø˜dÑ"�
Ø�MŠM˜*Ñ%Ô%Ð%Ð%Ø—K’K Ñ'Ô'ˆ	ØÐrC   c                óz  — d|› d| j         › �}| j                             |¦  «        }d                     t	          t          d„ |j        j        ¦  «        ¦  «        ¦  «        }	 |                      |¦  «        }d                     d„ |D ¦   «         ¦  «        }n# t          $ r d}Y nw xY w| j         › d|› d	|› d|› �S )
NzSELECT * FROM z LIMIT ú	c                ó   — | j         S rG   )Úname)Úfs    r2   rI   z1SparkSQL._get_sample_spark_rows.<locals>.<lambda>„   s   € °1´6€ rC   r[   c                ó8   — g | ]}d                       |¦  «        ‘ŒS )rg   )r_   )Ú.0rH   s     r2   ú
<listcomp>z3SparkSQL._get_sample_spark_rows.<locals>.<listcomp>ˆ   s"   € Ð(OÐ(OÐ(O¸C¨¯ª°3©¬Ð(OÐ(OÐ(OrC   Ú z rows from z table:
)
r.   r!   rJ   r_   rM   rN   r   ÚfieldsÚ_get_dataframe_resultsÚ	Exception)r/   rP   ÚqueryÚdfÚcolumns_strÚsample_rowsÚsample_rows_strs          r2   r]   zSparkSQL._get_sample_spark_rows�   sô   € ØP ÐPÐP¨tÔ/NÐPÐPˆØŒ[�_Š_˜UÑ#Ô#ˆØ—i’i¥¥SÐ)9Ð)9¸2¼9Ô;KÑ%LÔ%LÑ MÔ MÑNÔNˆð	!Ø×5Ò5°bÑ9Ô9ˆKà"ŸišiÐ(OÐ(OÀ;Ð(OÑ(OÔ(OÑPÔPˆOˆOøÝð 	!ð 	!ð 	!Ø ˆOˆOˆOð	!øøøð Ô.ð !ð !¸5ð !ð !Øð!ð !àð!ð !ð	
s   Á$4B ÂB(Â'B(rH   r
   Útuplec                óŽ   — t          t          t          |                     ¦   «                              ¦   «         ¦  «        ¦  «        S rG   )rw   rN   r6   ÚasDictÚvalues)r/   rH   s     r2   Ú_convert_row_as_tuplezSparkSQL._convert_row_as_tuple’   s.   € Ý•S�˜cŸjšj™lœl×1Ò1Ñ3Ô3Ñ4Ô4Ñ5Ô5Ð5rC   rs   r	   rM   c                ój   — t          t          | j        |                     ¦   «         ¦  «        ¦  «        S rG   )rM   rN   r{   rL   )r/   rs   s     r2   rp   zSparkSQL._get_dataframe_results•   s%   € Ý•C˜Ô2°B·J²J±L´LÑAÔAÑBÔBÐBrC   ÚallÚcommandÚfetchc                ó°   — | j                              |¦  «        }|dk    r|                     d¦  «        }t          |                      |¦  «        ¦  «        S )NÚoneé   )r!   rJ   Úlimitr6   rp   )r/   r~   r   rs   s       r2   ÚrunzSparkSQL.run˜   sI   € ØŒ[�_Š_˜WÑ%Ô%ˆØ�EŠ>ˆ>Ø—’˜!‘”ˆBÝ�4×.Ò.¨rÑ2Ô2Ñ3Ô3Ð3rC   c                óh   — 	 |                       |¦  «        S # t          $ r}	 d|› �cY d}~S d}~ww xY w)af  Get information about specified tables.

        Follows best practices as specified in: Rajkumar et al, 2022
        (https://arxiv.org/abs/2204.00498)

        If `sample_rows_in_table_info`, the specified number of sample rows will be
        appended to each table description. This can increase performance as
        demonstrated in the paper.
        úError: N)re   r(   )r/   rY   Úes      r2   Úget_table_info_no_throwz SparkSQL.get_table_info_no_throwž   sX   € ð	!Ø×&Ò& {Ñ3Ô3Ð3øÝð 	!ð 	!ð 	!Ø*Ø ˜Q�=�=Ð Ð Ð Ð Ð Ð øøøøð	!øøøs   ‚ —
1¡,¦1¬1c                ój   — 	 |                       ||¦  «        S # t          $ r}	 d|› �cY d}~S d}~ww xY w)a*  Execute a SQL command and return a string representing the results.

        If the statement returns rows, a string of the results is returned.
        If the statement returns no rows, an empty string is returned.

        If the statement throws an error, the error message is returned.
        r†   N)r„   rq   )r/   r~   r   r‡   s       r2   Úrun_no_throwzSparkSQL.run_no_throw®   sX   € ð	!Ø—8’8˜G UÑ+Ô+Ð+øÝð 	!ð 	!ð 	!Ø*Ø ˜Q�=�=Ð Ð Ð Ð Ð Ð øøøøð	!øøøs   ‚ ˜
2¢-§2­2)NNNNNr   )r   r   r   r   r   r   r   r   r   r   r   r   rG   )r5   r6   r7   r8   r9   r   r:   r   )r:   r@   )rP   r6   r:   r6   )rY   r   r:   r6   )rH   r
   r:   rw   )rs   r	   r:   rM   )r}   )r~   r6   r   r6   r:   r6   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r3   Úclassmethodr?   r*   r%   rX   re   r]   r{   rp   r„   rˆ   rŠ   © rC   r2   r   r   	   sU  € € € € € ØEÐEð 15Ø!%Ø $Ø-1Ø.2Ø)*ð=Dð =Dð =Dð =Dð =Dð~ à>Bð$ð $ð $ð $ñ „[ð$ð >ð >ð >ð >ð:ð :ð :ð :ð4ð 4ð 4ð 4ðð ð ð ð ð$
ð 
ð 
ð 
ð"6ð 6ð 6ð 6ðCð Cð Cð Cð4ð 4ð 4ð 4ð 4ð!ð !ð !ð !ð !ð !ð !ð !ð !ð !ð !ð !rC   r   N)Ú
__future__r   Útypingr   r   r   r   r   r   r	   r
   r   r   r�   rC   r2   ú<module>r“      s£   ðØ "Ð "Ð "Ð "Ð "Ð "à ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?Ð ?àð 9Ø8Ð8Ð8Ð8Ð8Ð8Ð8Ð8Ð8Ð8ðq!ð q!ð q!ð q!ð q!ñ q!ô q!ð q!ð q!ð q!rC   