Ë
    	êñi¥/  ã                  ó@  — d dl mZ d dlZd dlmZ d dlmZ d dlmZ d dl	m
Z
mZ d dlZd dlmZ d dlmZ d d	lmZ d d
lmZ d dlmZ d dlmZ d dlmZ e
r&d dlmZ d dlmZ d dlm Z m!Z! d dl"m#Z# d dl$m%Z% d dl&m'Z'  ed¬«       G d„ d«      «       Z(	 	 	 	 	 	 dd„Z)dd„Z*y)é    )ÚannotationsN)Ú	dataclass)Úpartial)Úperf_counter)ÚTYPE_CHECKINGÚAny)Úeprint)Úparse_version)Ú+_get_credentials_from_provider_expiry_aware)Ú0_extract_table_statistics_from_delta_add_actions)Úscan_parquet)ÚScanCastOptions)ÚSchema)Údatetime)Ú
DeltaTable)ÚDeletionFilesÚStorageOptionsDict)ÚNoPickleOption)ÚCredentialProviderBuilder©Ú	LazyFrameT)Úkw_onlyc                  óÀ   — e Zd ZU dZded<   ded<   ded<   ded	<   d
ed<   ded<   ded<   ded<   ded<   dd„Zddddddœ	 	 	 	 	 	 	 	 	 	 	 dd„Zdd„Zdd„Zdd„Z	dd„Z
y) ÚDeltaDatasetzDataset interface for Delta.zNoPickleOption[DeltaTable]Útable_ú
str | NoneÚ
table_uri_zint | str | datetime | NoneÚversionzStorageOptionsDict | NoneÚstorage_optionsz CredentialProviderBuilder | NoneÚcredential_provider_builderzdict[str, Any] | NoneÚdelta_table_optionsÚboolÚuse_pyarrowÚpyarrow_optionsÚrechunkc                óP   — t        | j                  «       j                  «       «      S )zFetch the schema of the table.)r   ÚtableÚschema©Úselfs    úZ/var/www/pod-logistic/pod-ai/venv/lib/python3.12/site-packages/polars/io/delta/_dataset.pyr(   zDeltaDataset.schema4   s   € ä�d—j‘j“l×)Ñ)Ó+Ó,Ð,ó    N)Úexisting_resolved_version_keyÚlimitÚ
projectionÚfilter_columnsÚpyarrow_predicatec               óÌ  ‡ ‡!— ddl Š ddl}|j                  j                  j	                  «       }|r.t        d| j                  › d|› d|› d|› d| j                  › �
«       | j                  «       Š!| j                  �| j                  n‰!j                  «       }t        |«      }	|�||	k(  r|rt        d|	›d	�«       y| j                  r„ddl
}dd
lm}
  ‰!j                  d!i | j                  xs i ¤Ž}t        |j                   j"                  j$                  j&                  ||||¬«      } |
j(                  |j*                  |dd¬«      |	fS ‰!j-                  «       }t/        |j0                  «      }| j+                  «       }t3        |j5                  «       D ��ci c]  \  }}||v sŒ||“Œ c}}«      }t7        «       }|rt        d«       ‰!j9                  «       }| j;                  «       j=                  d«      r|D �cg c]  }|j?                  dd«      ‘Œ }}|r)t7        «       |z
  }t        dtA        |«      › d|d›d�«       |�-tC         ‰ jD                  ‰!jG                  «       «      |||¬«      nd}‰!jI                  «       jJ                  }|duxr d|v }d}|rZddl&}d}tO        |jP                  «      }||k  r*ddjS                  d„ |D «       «      › d|› d�}tU        |«      ‚	 	 	 	 d"ˆ ˆ!fd„}d|f}nd}tW        |tA        |«      dkD  r|ndtA        |«      dkD  tY        jZ                  «       dd| j\                  | j^                  | j`                  ||¬ «      |	fS c c}}w c c}w )#zConstruct a LazyFrame scan.r   Nz*DeltaDataset: to_dataset_scan(): version: z	, limit: z, projection: z, filter_columns: z, use_pyarrow: z=DeltaDataset: to_dataset_scan(): early return (version_key = ú)r   )Ún_rowsÚ	predicateÚwith_columnsT)ÚpyarrowÚis_purez5DeltaDataset: to_dataset_scan(): begin path expansionz	lakefs://ús3://zCDeltaDataset: to_dataset_scan(): native scan_parquet(): num_files: z, path expansion time: z.3fÚs)r0   r(   ÚverboseÚdeletionVectors)é   é   é   z5reading delta deletion vectors requires deltalake >= ú.c              3  ó2   K  — | ]  }t        |«      –— Œ y ­w©N)Ústr)Ú.0Úvs     r+   ú	<genexpr>z/DeltaDataset.to_dataset_scan.<locals>.<genexpr>¦   s   è ø€ Ò,L¸¬S°¯VÑ,Lùs   ‚z, found c                ó´   •— t        ‰«      }|€? ‰j                  dd gt        | «      z  id ‰j                  ‰j                  «      i¬«      S t        | |«      S )NÚselection_vector)r(   )Ú_fetch_deletion_vectorsÚ	DataFrameÚlenÚListÚBooleanÚ_extract_delta_deletion_vectors)Úrequested_pathsÚdelta_deletion_vectorsÚplr'   s     €€r+   Ú_deletion_vector_callbackz?DeltaDataset.to_dataset_scan.<locals>._deletion_vector_callback«   sj   ø€ ô *AÀÓ)GÐ&Ø)Ð1Ø'˜2Ÿ<™<Ø+¨d¨V´c¸/Ó6JÑ-JÐKØ 2°G°B·G±G¸B¿J¹JÓ4GÐHôð ô 7Ø#Ð%;óð r,   zdelta-deletion-vectorÚinsertÚignore)
Úhive_schemaÚhive_partitioningÚcast_optionsÚmissing_columnsÚextra_columnsr   Úcredential_providerr%   Ú_table_statisticsÚ_deletion_files© )rO   úpl.DataFrameÚreturnr^   )1ÚpolarsÚpolars._utils.loggingÚ_utilsÚloggingr;   r	   r   r#   r'   rC   Ú(polars.io.pyarrow_dataset.anonymous_scanÚpolars.lazyframe.framer   Úto_pyarrow_datasetr$   r   ÚioÚpyarrow_datasetÚanonymous_scanÚ_scan_pyarrow_dataset_implÚ_scan_python_functionr(   ÚmetadataÚsetÚpartition_columnsr   Úitemsr   Ú	file_urisÚ	table_uriÚ
startswithÚreplacerK   r   rJ   Úget_add_actionsÚprotocolÚreader_featuresÚ	deltalaker
   Ú__version__ÚjoinÚImportErrorr   r   Ú_default_icebergr   r    r%   )"r*   r-   r.   r/   r0   r1   r`   r;   r   Úversion_keyr   ÚdatasetÚfuncÚtable_mdrn   r(   ÚkrE   rU   Ú
start_timeÚpathsÚpathÚelapsedÚtable_statisticsrv   Úhas_deletion_vectorsÚdeletion_filesrw   Údv_min_versionÚ	installedÚmsgrR   rQ   r'   s"                                   @@r+   Úto_dataset_scanzDeltaDataset.to_dataset_scan8   s°  ù€ ó 	Û$à—-‘-×'Ñ'×/Ñ/Ó1ˆáÜðØ ŸL™L˜>ð *Ø˜ð !Ø)˜lð +#Ø#1Ð"2ð 3 Ø $× 0Ñ 0Ð1ð3ôð —
‘
“ˆØ"&§,¡,Ð":�$—,’,ÀÇÁÃˆÜ˜'“lˆð *Ð5Ø-°Ò<áÜØTÀkÐEUÐUVÐWôð à×ÒÛ;Ý8à.�e×.Ñ.ÑN°$×2FÑ2FÒ2LÈ"ÑNˆGäØ—	‘	×)Ñ)×8Ñ8×SÑSØØØ+Ø'ôˆDð 3�9×2Ñ2Ø—‘ ¨d¸Dôàðð ð —>‘>Ó#ˆÜ × :Ñ :Ó;Ðà—‘“ˆÜØ$Ÿl™l›n×G‘d�a˜°Ð5FÒ0FˆQ�‰TÓGó
ˆô "“^ˆ
áÜÐJÔKà—‘Ó!ˆà�>‰>Ó×&Ñ& {Ô3ØDIÖJ¸D�T—\‘\ +¨wÕ7ÐJˆEÐJáÜ"“n zÑ1ˆGÜðä! %›j˜\ð *(Ø(/° }°Að7ôð Ð)ô =Ø�—‘˜U×2Ñ2Ó4Ó5Ø-ØØõ	ð ð 	ð  Ÿ.™.Ó*×:Ñ:ˆà 4Ð'ÒPÐ,=ÀÐ,Pð 	ð 04ˆÙÛà&ˆNÜ% i×&;Ñ&;Ó<ˆIØ˜>Ò)ð$Ø$'§H¡HÑ,L¸^Ô,LÓ$LÐ#Mð NØ&˜K qð*ð ô
 " #Ó&Ð&ðØ!-ðàöð (Ø)ð‰Nð
 "ˆNäØÜ'*Ð+<Ó'=ÀÒ'A™ÀtÜ!Ð"3Ó4°qÑ8Ü(×9Ñ9Ó;Ø$Ø"Ø ×0Ñ0Ø $× @Ñ @Ø—L‘LØ.Ø*ô
ð ðð 	ùóQ Hùò Ks   ÆM
Æ%M
Ç<M!c                ó¨   — | j                   €;| j                  j                  «       €J ‚| j                  «       j                  | _         | j                   S )zFetch the table URI.)r   r   Úgetr'   rq   r)   s    r+   rq   zDeltaDataset.table_uriÑ   s@   € à�?‰?Ð"Ø—;‘;—?‘?Ó$Ð0Ð0Ð0Ø"Ÿj™j›l×4Ñ4ˆDŒOà�‰Ðr,   c                óh  — | j                   j                  «       �€~ddlm} ddlm}m}m} |j                  d«       ddl	m
} | j                  €J ‚i }| j                  r+| j                  j                  «       x}rt        |«      xs i } || j                  | j                  | j                   €| j                  �i | j                   xs i ¥|¥nd| j"                  ¬«      }|j%                  «       }	|	j&                  |kD  s|	j&                  |k(  rd|	j&                  › d	|› d
|› �}
 ||
«      ‚|	j&                  dk\  rE|	j(                  �9h |	j(                  £j+                  |«      }t-        |«      dkD  rd|› d�}
 ||
«      ‚| j                   j/                  |«       | j                   j                  «       S )zFetch the DeltaTable object.Nr   )ÚDeltaProtocolError)ÚMAX_SUPPORTED_READER_VERSIONÚNOT_SUPPORTED_READER_VERSIONÚSUPPORTED_READER_FEATURESr<   )Ú_get_delta_lake_table)Ú
table_pathr   r   r!   z&The table's minimum reader version is z5 but polars delta scanner only supports version 1 or z with these reader features: é   z)The table has set these reader features: z= but these are not yet supported by the polars delta scanner.)r   r�   Údeltalake.exceptionsr�   Údeltalake.tabler�   r‘   r’   ÚaddÚpolars.io.delta._utilsr“   r   r    Úbuild_credential_providerr   r   r   r!   ru   Úmin_reader_versionrv   Ú
differencerK   rm   )r*   r�   r�   r‘   r’   r“   Úcredential_provider_credsÚproviderr'   Útable_protocolrŠ   Úmissing_featuress               r+   r'   zDeltaDataset.tableÙ   så  € à�;‰;�?‰?ÓÑ$Ý?÷ñ ð &×)Ñ)Ð*;Ô<åDà—?‘?Ð.Ð.Ð.à(*Ð%à×/Ò/Ø ×<Ñ<×VÑVÓXÐX�ÐXô @ÀÓIÒOÈRð *ñ *ØŸ?™?ØŸ™ð ×+Ñ+Ð7Ø×7Ñ7ÐCð R˜×,Ñ,Ò2°ÐQÐ7PÑQð à$(×$<Ñ$<ô
ˆEð #Ÿ^™^Ó-ˆNð ×1Ñ1Ð4PÒPØ!×4Ñ4Ð8TÒTð =¸^×=^Ñ=^Ð<_ð `KØKgÐJhð  iFð  G`ð  Faðbð ñ )¨Ó-Ð-à×1Ñ1°QÒ6Ø"×2Ñ2Ð>à#D ^×%CÑ%CÐ#D×#OÑ#OØ-ó$Ð ô Ð'Ó(¨1Ò,ØEÐFVÐEWð  XUð  V�CÙ,¨SÓ1Ð1à�K‰K�O‰O˜EÔ"à�{‰{�‰Ó Ð r,   c                ó:   — | j                  «        | j                  S rB   )rq   Ú__dict__r)   s    r+   Ú__getstate__zDeltaDataset.__getstate__  s   € Ø�‰ÔØ�}‰}Ðr,   c                ó   — || _         y rB   )r¢   )r*   Ústates     r+   Ú__setstate__zDeltaDataset.__setstate__  s	   € Øˆ�r,   )r_   r   )r-   r   r.   z
int | Noner/   úlist[str] | Noner0   r§   r1   r   r_   ztuple[LazyFrame, str] | None)r_   rC   )r_   r   )r_   údict[str, Any])r¥   r¨   r_   ÚNone)Ú__name__Ú
__module__Ú__qualname__Ú__doc__Ú__annotations__r(   r‹   rq   r'   r£   r¦   r]   r,   r+   r   r      s®   … á&à&Ó&ØÓØ(Ó(à.Ó.Ø!AÓAØ.Ó.àÓØ*Ó*àƒMó-ð 59Ø Ø'+Ø+/Ø(,ñSð (2ðSð ð	Sð
 %ðSð )ðSð &ðSð 
&óSóró>!ó@ôr,   r   c                óf  — | j                   dt        j                  ik(  sJ ‚t        j                  t        j                  t        j                  «      dœ}|j                  |j                  «       «      }|j                   |k(  sJ ‚t        j                  dk7  rdnd}| j                  «       j                  t        j                  d«      j                  j                  dd«      j                  j                  |«      «      j                  |j                  «       j                  t        j                  d«      j                  j                  dd«      j                  j                  |«      «      ddd	d	¬
«      j                  dg«      j!                  «       }|j"                  t%        | «      k(  sJ ‚|S )a  
    Extract the deletion_vectors for the provided requested_paths.

    Input requested_paths schema is "path": String.
    Output series schema is "selection_vector": List(Boolean), maintaining order.

    The selection_vector from deltalake is a keep-mask (True = keep).
    rƒ   )ÚfilepathrH   Úwin32zfile://zfile:///z
^lakefs://r9   r°   Úleft)Úleft_onÚright_onÚhowÚmaintain_orderrH   )r(   rQ   ÚStringrL   rM   ÚselectÚkeysÚsysÚplatformÚlazyr6   ÚcolrC   rs   Ústrip_prefixry   ÚcollectÚheightrK   )rO   rP   Údelta_dv_schemaÚfile_prefixÚ	joined_dfs        r+   rN   rN   !  sX  € ð ×!Ñ! f¬b¯i©iÐ%8Ò8Ð8Ð8ä#%§9¡9Ä"Ç'Á'Ì"Ï*É*ÓBUÑV€OØ3×:Ñ:¸?×;OÑ;OÓ;QÓRÐØ!×(Ñ(¨OÒ;Ð;Ð;ä"Ÿ|™|¨wÒ6‘)¸J€Kà×ÑÓß	‰Ü�F‰F�6‹Nß‰S—‘˜ wÓ/ß‰S—‘˜kÓ*ó

÷
 
‰Ø"×'Ñ'Ó)×6Ñ6Ü—‘�zÓ"ß‘—W‘W˜\¨7Ó3ß‘—\‘\ +Ó.óð
 ØØØ!ð 
ó 


÷ 
‰Ð#Ð$Ó	%ß	‰‹ð' ð, ×Ñœs ?Ó3Ò3Ð3Ð3àÐr,   c                ó
  — ddl }|j                  j                  j                  «       }t	        j
                  | j                  «       «      }|r&|j                  dkD  rt        dt        |«      › �«       t        |«      dk(  ry|S )aT  
    Fetch the deletion_vectors, mapping file_uri to "deletion_vector".

    Schema: {"filepath": pl.String, "selection_vector": pl.List(pl.Boolean)}

    The selection_vector from deltalake is a keep-mask (True = keep), so
    the more accurate term would be "selection_vector".

    Returns None if the table has no deletion vectors.
    r   Nz0DeltaDataset: has deletion_vectors, file_count: )
ra   rb   rc   r;   rQ   rJ   Údeletion_vectorsrÀ   r	   rK   )r'   r`   r;   Údv_tables       r+   rI   rI   O  sl   € ó !à�m‰m×#Ñ#×+Ñ+Ó-€Gä�|‰|˜E×2Ñ2Ó4Ó5€Há�8—?‘? QÒ&ÜÐAÄ#ÀhÃ-ÀÐQÔRä
ˆ8ƒ}˜ÒØà€Or,   )rO   r^   rP   r^   r_   r^   )r'   r   r_   zpl.DataFrame | None)+Ú
__future__r   rº   Údataclassesr   Ú	functoolsr   Útimer   Útypingr   r   r`   rQ   ra   r	   Úpolars._utils.variousr
   Ú.polars.io.cloud.credential_provider._providersr   r™   r   Úpolars.io.parquet.functionsr   Ú#polars.io.scan_options.cast_optionsr   Úpolars.schemar   r   rw   r   Úpolars._typingr   r   Úpolars.io.cloud._utilsr   Ú,polars.io.cloud.credential_provider._builderr   re   r   r   rN   rI   r]   r,   r+   ú<module>rÔ      s�   ðÝ "ã 
Ý !Ý Ý ß %ã Ý (Ý /õõ TÝ 4Ý ?Ý  áÝ!å$ç@Ý5ÝVÝ0ñ �4Ô÷~ð ~ó ð~ðB+Ø!ð+à(ð+ð ó+ô\r,   