U ŸDx`Ò-ã@s„ddlmZddlmZddlmZmZddlZddlm Z ddl Z ddl Z ddl Z ddlZddlZddlZddlmZddlmZddlmZmZmZmZmZmZmZddlmZm Z m!Z!m"Z"ddlm#Z$dd l%m&Z&m'Z'm(Z(d Z)d d „Z*d d„Z+dd„Z,d]dd„Z-dZ.dd„Z/Gdd„dƒZ0e  1d¡Z2dd„Z3dd„Z4dd„Z5d Z6Gd!d"„d"ƒZ7d#d$„Z8Gd%d&„d&ƒZ9Gd'd(„d(ƒZ:Gd)d*„d*ƒZ;Gd+d,„d,ƒZd1d2„Z?d3hZ@Gd4d5„d5ƒZAd^d6d7„ZBd8ZCGd9d:„d:ƒZDd_d=d>„ZEGd?d@„d@ƒZFdAZGd`dDdE„ZHeG IdFdG JeCdHf¡dIe.¡eH_KdadJdK„ZLeG IdLeCdMe.¡eL_KdbdPdQ„ZMdR Ie6¡eM_KdSdT„ZNdcdUdV„ZOdddWdX„ZPdedYdZ„ZQdfd[d\„ZRdS)gé)Ú defaultdict)Úfutures)ÚpartialÚreduceN)Ú Collection)Ú ParquetReaderÚ StatisticsÚ FileMetaDataÚRowGroupMetaDataÚColumnChunkMetaDataÚ ParquetSchemaÚ ColumnSchema)ÚLocalFileSystemÚ FileSystemÚ_resolve_filesystem_and_pathÚ_ensure_filesystem)Ú filesystem)ÚguidÚ _is_path_likeÚ_stringify_path)ZhdfscCs,t|ƒ}tj |¡}|jtkr$|jS|SdS©N)rÚurllibÚparseÚurlparseÚschemeÚ_URI_STRIP_SCHEMESÚpath)rZ parsed_uri©rú6/tmp/pip-target-oguziej0/lib/python/pyarrow/parquet.pyÚ _parse_uri/s   rcCs2|dkrt ||¡St |¡}t|ƒ}||fSdSr)ÚlegacyfsÚresolve_filesystem_and_pathrr)Zpassed_filesystemrZ parsed_pathrrrÚ_get_filesystem_and_path:s   r"cCsRt|tƒr<|D]*}t|tƒr&tdƒ}nd}||krdSqnt|tƒrNd|kSdS)NrTúF)Ú isinstanceÚbytesÚchrÚstr)ÚvalÚbyteZ compare_torrrÚ_check_contains_nullCs     r*TcCs”|dk r�t|ƒdks&tdd„|Dƒƒr.tdƒ‚t|ddtƒrF|g}|r�|D]@}|D]6\}}}t|tƒr|tdd„|Dƒƒs„t|ƒrVtdƒ‚qVqN|S)z+ Check if filters are well-formed. Nrcss|]}t|ƒdkVqdS)rN)Úlen©Ú.0ÚfrrrÚ Vsz!_check_filters..zMalformed filterscss|]}t|ƒVqdSr)r*)r-Úvrrrr/bszBNull-terminated binary strings are not supported as filter values.) r+ÚanyÚ ValueErrorr$r'ÚlistÚallr*ÚNotImplementedError)ÚfiltersÚcheck_null_stringsÚ conjunctionÚcolÚopr(rrrÚ_check_filtersQs$ÿþýÿr;azPredicates are expressed in disjunctive normal form (DNF), like ``[[('x', '=', 0), ...], ...]``. DNF allows arbitrary boolean logical combinations of single column predicates. The innermost tuples each describe a single column predicate. The list of inner predicates is interpreted as a conjunction (AND), forming a more selective and multiple column predicate. Finally, the most outer list combines these filters as a disjunction (OR). Predicates may also be passed as List[Tuple]. This form is interpreted as a single conjunction. To express OR in predicates, one must use the (preferred) List[List[Tuple]] notation. Each tuple has format: (``key``, ``op``, ``value``) and compares the ``key`` with the ``value``. The supported ``op`` are: ``=`` or ``==``, ``!=``, ``<``, ``>``, ``<=``, ``>=``, ``in`` and ``not in``. If the ``op`` is ``in`` or ``not in``, the ``value`` must be a collection such as a ``list``, a ``set`` or a ``tuple``. Examples: .. code-block:: python ('x', '=', 0) ('y', 'in', ['a', 'b', 'c']) ('z', 'not in', {'a','b'}) csrddlm‰t|ˆjƒr|St|dd�}‡fdd„‰g}|D](}‡fdd„|Dƒ}| ttj|ƒ¡qú<=ú>=Úinúnot inz,"{0}" is not a valid operator in predicates.)ÚfieldÚisinr2Úformat)r9r:r(rE)ÚdsrrÚconvert_single_predicate—s,   ÿÿz8_filters_to_expression..convert_single_predicatecsg|]\}}}ˆ|||ƒ‘qSrr)r-r9r:r()rIrrÚ ²sÿz*_filters_to_expression..) Úpyarrow.datasetÚdatasetr$Z Expressionr;ÚappendrÚoperatorÚand_Úor_)r6Zdisjunction_membersr8Zconjunction_membersr)rIrHrÚ_filters_to_expressionŠs     þrQc@sŽeZdZdZddd„Zdd„Zed d „ƒZed d „ƒZed d„ƒZ edd„ƒZ d dd„Z d!dd„Z d"dd„Z d#dd„Zd$dd„Zd%dd„ZdS)&Ú ParquetFilea‹ Reader interface for a single Parquet file. Parameters ---------- source : str, pathlib.Path, pyarrow.NativeFile, or file-like object Readable source. For passing bytes or buffer-like file containing a Parquet file, use pyarrow.BufferReader. metadata : FileMetaData, default None Use existing metadata object, rather than reading from file. common_metadata : FileMetaData, default None Will be used in reads for pandas schema metadata if not found in the main file's metadata, no other uses at the moment. memory_map : bool, default False If the source is a file path, use a memory map to read file, which can improve performance in some environments. buffer_size : int, default 0 If positive, perform read buffering when deserializing individual column chunks. Otherwise IO calls are unbuffered. NFrcCs2tƒ|_|jj|||||d�||_| ¡|_dS)N)Zuse_memory_mapÚ buffer_sizeÚread_dictionaryÚmetadata)rÚreaderÚopenÚcommon_metadataÚ_build_nested_pathsÚ_nested_paths_by_prefix)ÚselfÚsourcerUrXrTÚ memory_maprSrrrÚ__init__Ös þzParquetFile.__init__cCsn|jj}ttƒ}t|ƒD]P\}}|d}|dd…}|| |¡|sHqd ||df¡}|dd…}q4q|S)NréÚ.)rVZ column_pathsrr3Ú enumeraterMÚjoin)r[ÚpathsÚresultÚirÚkeyÚrestrrrrYßs zParquetFile._build_nested_pathscCs|jjSr)rVrU©r[rrrrUòszParquetFile.metadatacCs|jjS)zG Return the Parquet schema, unconverted to Arrow types )rUÚschemarhrrrriöszParquetFile.schemacCs|jjS)zj Return the inferred Arrow schema, converted from the whole Parquet file's schema )rVÚ schema_arrowrhrrrrjýszParquetFile.schema_arrowcCs|jjSr)rVÚnum_row_groupsrhrrrrkszParquetFile.num_row_groupsTcCs |j||d�}|jj|||d�S)a¾ Read a single row group from a Parquet file. Parameters ---------- columns: list If not None, only these columns will be read from the row group. A column name may be a prefix of a nested field, e.g. 'a' will select 'a.b', 'a.c', and 'a.d.e'. use_threads : bool, default True Perform multi-threaded column reads. use_pandas_metadata : bool, default False If True and file has custom pandas schema metadata, ensure that index columns are also loaded. Returns ------- pyarrow.table.Table Content of the row group as a table (of columns) ©Úuse_pandas_metadata©Úcolumn_indicesÚ use_threads)Ú_get_column_indicesrVÚread_row_group)r[reÚcolumnsrprmrorrrrr sÿ ÿzParquetFile.read_row_groupcCs |j||d�}|jj|||d�S)a Read a multiple row groups from a Parquet file. Parameters ---------- row_groups: list Only these row groups will be read from the file. columns: list If not None, only these columns will be read from the row group. A column name may be a prefix of a nested field, e.g. 'a' will select 'a.b', 'a.c', and 'a.d.e'. use_threads : bool, default True Perform multi-threaded column reads. use_pandas_metadata : bool, default False If True and file has custom pandas schema metadata, ensure that index columns are also loaded. Returns ------- pyarrow.table.Table Content of the row groups as a table (of columns). rlrn)rqrVÚread_row_groups)r[Ú row_groupsrsrprmrorrrrt$sÿþzParquetFile.read_row_groupsécCs<|dkrtd|jjƒ}|j||d�}|jj||||d�}|S)aà Read streaming batches from a Parquet file Parameters ---------- batch_size: int, default 64K Maximum number of records to yield per batch. Batches may be smaller if there aren't enough rows in the file. row_groups: list Only these row groups will be read from the file. columns: list If not None, only these columns will be read from the file. A column name may be a prefix of a nested field, e.g. 'a' will select 'a.b', 'a.c', and 'a.d.e'. use_threads : boolean, default True Perform multi-threaded column reads. use_pandas_metadata : boolean, default False If True and file has custom pandas schema metadata, ensure that index columns are also loaded. Returns ------- iterator of pyarrow.RecordBatch Contents of each batch as a record batch Nrrl)rurorp)ÚrangerUrkrqrVÚ iter_batches)r[Ú batch_sizerursrprmroZbatchesrrrrxBsÿýzParquetFile.iter_batchescCs|j||d�}|jj||d�S)aª Read a Table from Parquet format, Parameters ---------- columns: list If not None, only these columns will be read from the file. A column name may be a prefix of a nested field, e.g. 'a' will select 'a.b', 'a.c', and 'a.d.e'. use_threads : bool, default True Perform multi-threaded column reads. use_pandas_metadata : bool, default False If True and file has custom pandas schema metadata, ensure that index columns are also loaded. Returns ------- pyarrow.table.Table Content of the file as a table (of columns). rlrn)rqrVZread_all)r[rsrprmrorrrÚreadhsÿÿzParquetFile.readcCs| |¡}|jj||d�S)a Read contents of file for the given columns and batch size. Notes ----- This function's primary purpose is benchmarking. The scan is executed on a single thread. Parameters ---------- columns : list of integers, default None Select columns to read, if None scan all columns. batch_size : int, default 64K Number of rows to read at a time internally. Returns ------- num_rows : number of rows in file )ry)rqrVÚ scan_contents)r[rsryrorrrr{‚s ÿzParquetFile.scan_contentscs¬|dkr dSg}|D]}|ˆjkr| ˆj|¡q|r¨ˆjj}ˆjdk rRˆjjnd}|rld|krlt|ƒ}n|r‚d|kr‚t|ƒ}ng}|dk r¨|r¨|‡fdd„|Dƒ7}|S)Nópandascs"g|]}t|tƒsˆj |¡‘qSr)r$ÚdictrVZcolumn_name_idx)r-ÚdescrrhrrrJ²s þz3ParquetFile._get_column_indices..)rZÚextendrUrXÚ_get_pandas_index_columns)r[Z column_namesrmÚindicesÚnameZfile_keyvaluesZcommon_keyvaluesÚ index_columnsrrhrrqšs, ÿ þ      ÿzParquetFile._get_column_indices)NNNFr)NTF)NTF)rvNNTF)NTF)Nrv)F)Ú__name__Ú __module__Ú __qualname__Ú__doc__r^rYÚpropertyrUrirjrkrrrtrxrzr{rqrrrrrRÀs8ÿ     ÿ ÿ ÿ &  rRz [ ,;{}() =]cCs t d|¡S)NÚ_)Ú_SPARK_DISALLOWED_CHARSÚsub)r‚rrrÚ_sanitized_spark_field_name¼srŒc Cs„d|krxg}d}|D]J}|j}t|ƒ}||krTd}t ||j|j|j¡}| |¡q| |¡qtj||jd�}||fS|dfSdS)NÚsparkFT)rU) r‚rŒÚparEÚtypeZnullablerUrMri) riÚflavorZsanitized_fieldsÚschema_changedrEr‚Zsanitized_nameZsanitized_fieldÚ new_schemarrrÚ_sanitize_schemaÀs" ÿ  r“cs8d|kr0‡fdd„tˆjƒDƒ}tjj||d�SˆSdS)Nr�csg|] }ˆ|‘qSrr)r-re©ÚtablerrrJÛsz#_sanitize_table..)ri)rwZ num_columnsrŽÚTableÚ from_arrays)r•r’r�Z column_datarr”rÚ_sanitize_tableØsr˜a version : {"1.0", "2.0"}, default "1.0" Determine which Parquet logical types are available for use, whether the reduced set from the Parquet 1.x.x format or the expanded logical types added in format version 2.0.0 and after. Note that files written with version='2.0' may not be readable in all Parquet implementations, so version='1.0' is likely the choice that maximizes file compatibility. Some features, such as lossless storage of nanosecond timestamps as INT64 physical storage, are only available with version='2.0'. The Parquet 2.0.0 format version also introduced a new serialized data page format; this can be enabled separately using the data_page_version option. use_dictionary : bool or list Specify if we should use dictionary encoding in general or only for some columns. use_deprecated_int96_timestamps : bool, default None Write timestamps to INT96 Parquet format. Defaults to False unless enabled by flavor argument. This take priority over the coerce_timestamps option. coerce_timestamps : str, default None Cast timestamps a particular resolution. The defaults depends on `version`. For ``version='1.0'`` (the default), nanoseconds will be cast to microseconds ('us'), and seconds to milliseconds ('ms') by default. For ``version='2.0'``, the original resolution is preserved and no casting is done by default. The casting might result in loss of data, in which case ``allow_truncated_timestamps=True`` can be used to suppress the raised exception. Valid values: {None, 'ms', 'us'} data_page_size : int, default None Set a target threshold for the approximate encoded size of data pages within a column chunk (in bytes). If None, use the default data page size of 1MByte. allow_truncated_timestamps : bool, default False Allow loss of data when coercing timestamps to a particular resolution. E.g. if microsecond or nanosecond data is lost when coercing to 'ms', do not raise an exception. compression : str or dict Specify the compression codec, either on a general basis or per-column. Valid values: {'NONE', 'SNAPPY', 'GZIP', 'BROTLI', 'LZ4', 'ZSTD'}. write_statistics : bool or list Specify if we should write statistics in general (default is True) or only for some columns. flavor : {'spark'}, default None Sanitize schema or set other compatibility options to work with various target systems. filesystem : FileSystem, default None If nothing passed, will be inferred from `where` if path-like, else `where` is already a file-like object so no filesystem is needed. compression_level: int or dict, default None Specify the compression level for a codec, either on a general basis or per-column. If None is passed, arrow selects the compression level for the compression codec in use. The compression level has a different meaning for each codec, so you have to read the documentation of the codec you are using. An exception is thrown if the compression codec does not allow specifying a compression level. use_byte_stream_split: bool or list, default False Specify if the byte_stream_split encoding should be used in general or only for some columns. If both dictionary and byte_stream_stream are enabled, then dictionary is preferred. The byte_stream_split encoding is valid only for floating-point data types and should be combined with a compression codec. data_page_version : {"1.0", "2.0"}, default "1.0" The serialized Parquet data page format version to write, defaults to 1.0. This does not impact the file schema logical types and Arrow to Parquet type casting behavior; for that use the "version" option. c @sJeZdZd e¡Zddd„Zd d „Zd d „Zd d„Z ddd„Z dd„Z dS)Ú ParquetWriteraˆ Class for incrementally building a Parquet file for Arrow tables. Parameters ---------- where : path or file-like object schema : arrow Schema {} **options : dict If options contains a key `metadata_collector` then the corresponding value is assumed to be a list (or any object with `.append` method) that will be filled with the file metadata instance of the written file. Nú1.0TÚsnappyFc Ksô| dkr"|dk rd|krd} nd} ||_|dk rBt||ƒ\}|_nd|_||_||_d|_t||dd�\}}|dk rªt|tj ƒr”|  |d¡}|_q®|j |dd�}|_n|}|  dd¡|_ d}tj||f||||| | | || d œ |—Ž|_d|_dS) Nr�TF)Zallow_legacy_filesystemÚwb)Ú compressionÚmetadata_collectorZV2) Úversionr�Úuse_dictionaryÚwrite_statisticsÚuse_deprecated_int96_timestampsÚcompression_levelÚuse_byte_stream_splitÚwriter_engine_versionÚdata_page_version)r�r“r‘riÚwhereÚ file_handlerr$r rrWZopen_output_streamÚpopÚ_metadata_collectorÚ_parquetr™ÚwriterÚis_open)r[r§rirr�rŸr r�r¡r¢r£r¤r¥r¦ÚoptionsrZsinkZengine_versionrrrr^4sV ÿ  ÿÿö õ zParquetWriter.__init__cCst|ddƒr| ¡dS)Nr­F)ÚgetattrÚcloserhrrrÚ__del__ts zParquetWriter.__del__cCs|SrrrhrrrÚ __enter__xszParquetWriter.__enter__cOs | ¡dS©NF)r°)r[ÚargsÚkwargsrrrÚ__exit__{szParquetWriter.__exit__cCs^|jrt||j|jƒ}|js t‚|jj|jdd�sJd |j|j¡}t|ƒ‚|j j ||d�dS)NF©Zcheck_metadatazTTable schema does not match schema used to create file: table: {!s} vs. file: {!s}©Úrow_group_size) r‘r˜rir�r­ÚAssertionErrorÚequalsrGr2r¬Ú write_table)r[r•r¹Úmsgrrrr¼€s þzParquetWriter.write_tablecCsH|jr0|j ¡d|_|jdk r0|j |jj¡|jdk rD|j ¡dSr³)r­r¬r°rªrMrUr¨rhrrrr°�s   zParquetWriter.close) NNršTr›TNNFNrš)N) r„r…r†rGÚ_parquet_writer_arg_docsr‡r^r±r²r¶r¼r°rrrrr™#s( óö @ r™cCst |d d¡¡dS)Nr|Úutf8rƒ)ÚjsonÚloadsÚdecode)Ú keyvaluesrrrr€—sÿr€c@s\eZdZdZeedd�dddfdd„Zdd„Zd d „Zd d „Z d d„Z dd„Zddd„Z dS)ÚParquetDatasetPiecea� A single chunk of a potentially larger Parquet dataset to read. The arguments will indicate to read either a single row group or all row groups, and whether to add partition keys to the resulting pyarrow.Table. Parameters ---------- path : str or pathlib.Path Path to file in the file system where this piece is located. open_file_func : callable Function to use for obtaining file handle to dataset piece. partition_keys : list of tuples Two-element tuples of ``(column name, ordinal index)``. row_group : int, default None Row group to load. By default, reads all row groups. Úrb©ÚmodeNcCs.t|ƒ|_||_||_|pg|_|p&i|_dSr)rrÚopen_file_funcÚ row_groupÚpartition_keysÚ file_options)r[rrÈrËrÉrÊrrrr^´s   zParquetDatasetPiece.__init__cCs2t|tƒsdS|j|jko0|j|jko0|j|jkSr³)r$rÄrrÉrÊ©r[ÚotherrrrÚ__eq__¼s   ÿ þzParquetDatasetPiece.__eq__cCsd t|ƒj|j|j|j¡S)Nz-{}({!r}, row_group={!r}, partition_keys={!r}))rGr�r„rrÉrÊrhrrrÚ__repr__Ãs ýzParquetDatasetPiece.__repr__cCs^d}t|jƒdkr6d dd„|jDƒ¡}|d |¡7}||j7}|jdk rZ|d |j¡7}|S)NÚrz, css|]\}}d ||¡VqdS)z{}={}N©rG)r-r‚Úindexrrrr/Ísÿz.ParquetDatasetPiece.__str__..zpartition[{}] z | row_group={})r+rÊrbrGrrÉ)r[rdZ partition_strrrrÚ__str__És ÿ  zParquetDatasetPiece.__str__cCs| ¡}|jS)zn Return the file's metadata. Returns ------- metadata : FileMetaData )rWrU)r[r.rrrÚ get_metadataØsz ParquetDatasetPiece.get_metadatacCs(| |j¡}t|tƒs$t|f|jŽ}|S)z1 Return instance of ParquetFile. )rÈrr$rRrË)r[rVrrrrWãs  zParquetDatasetPiece.openTFcCsæ|jdk r| ¡}n(|dk r,t|f|jŽ}nt|jf|jŽ}t|||d�}|jdk rf|j|jf|Ž}n |jf|Ž}t |j ƒdkrâ|dkr�t dƒ‚t |j ƒD]F\} \} } t jt |ƒ| dd�} |j| j} tj | | ¡}| | |¡}qš|S)a¢ Read this piece as a pyarrow.Table. Parameters ---------- columns : list of column names, default None use_threads : bool, default True Perform multi-threaded column reads. partitions : ParquetPartitions, default None file : file-like object Passed to ParquetFile. Returns ------- table : pyarrow.Table N©rsrprmrzMust pass partition setsÚi4)Zdtype)rÈrWrRrËrr}rÉrrrzr+rÊr2raÚnpÚfullÚlevelsÚ dictionaryrŽZDictionaryArrayr—Z append_column)r[rsrpÚ partitionsÚfilermrVr®r•rer‚rÒr�rÚZarrrrrrzìs*  þ    zParquetDatasetPiece.read)NTNNF) r„r…r†r‡rrWr^rÎrÏrÓrÔrzrrrrrÄ¡s ÿ   ÿrÄc@s:eZdZdZd dd„Zdd„Zedd„ƒZed d „ƒZdS) Ú PartitionSeta¼ A data structure for cataloguing the observed Parquet partitions at a particular level. So if we have /foo=a/bar=0 /foo=a/bar=1 /foo=a/bar=2 /foo=b/bar=0 /foo=b/bar=1 /foo=b/bar=2 Then we have two partition sets, one for foo, another for bar. As we visit levels of the partition hierarchy, a PartitionSet tracks the distinct values and assigns categorical codes to use when reading the pieces NcCs0||_|p g|_dd„t|jƒDƒ|_d|_dS)NcSsi|]\}}||“qSrr)r-reÚkrrrÚ @sz)PartitionSet.__init__..)r‚ÚkeysraÚ key_indicesÚ _dictionary)r[r‚ràrrrr^=s zPartitionSet.__init__cCs<||jkr|j|St|jƒ}|j |¡||j|<|SdS)zc Get the index of the partition value if it is known, otherwise assign one N)rár+ràrM)r[rfrÒrrrÚ get_indexCs      zPartitionSet.get_indexcCsp|jdk r|jSt|jƒdkr&tdƒ‚zdd„|jDƒ}t |¡}Wn tk rdt |j¡}YnX||_|S)NrzNo known partition keyscSsg|] }t|ƒ‘qSr)Úint©r-ÚxrrrrJZsz+PartitionSet.dictionary..)râr+ràr2ÚlibÚarray)r[Z integer_keysrÚrrrrÚPs zPartitionSet.dictionarycCst|jƒt|jƒkSr)r3ràÚsortedrhrrrÚ is_sortedbszPartitionSet.is_sorted)N) r„r…r†r‡r^rãrˆrÚrêrrrrrÝ,s   rÝc@sDeZdZdd„Zdd„Zdd„Zdd„Zd d „Zd d „Zd d„Z dS)ÚParquetPartitionscCsg|_tƒ|_dSr)rÙÚsetÚpartition_namesrhrrrr^iszParquetPartitions.__init__cCs t|jƒSr)r+rÙrhrrrÚ__len__mszParquetPartitions.__len__cCs |j|Sr)rÙ)r[rerrrÚ __getitem__pszParquetPartitions.__getitem__cCs*t|tƒstdƒ‚|j|jko(|j|jkS)Nz0`other` must be an instance of ParquetPartitions)r$rëÚ TypeErrorrÙrírÌrrrr»ss    ÿzParquetPartitions.equalscCs*z | |¡WStk r$tYSXdSr©r»rðÚNotImplementedrÌrrrrÎzs zParquetPartitions.__eq__cCsV|t|jƒkrF||jkr&td |¡ƒ‚t|ƒ}|j |¡|j |¡|j| |¡S)aT Record a partition value at a particular level, returning the distinct code for that value at that level. Example: partitions.get_index(1, 'foo', 'a') returns 0 partitions.get_index(1, 'foo', 'b') returns 1 partitions.get_index(1, 'foo', 'c') returns 2 partitions.get_index(1, 'foo', 'a') returns 0 Parameters ---------- level : int The nesting level of the partition we are observing name : str The partition name key : str or int The partition value z1{} was the name of the partition in another level) r+rÙrír2rGrÝrMÚaddrã)r[Úlevelr‚rfZpart_setrrrrã€s ÿ  zParquetPartitions.get_indexc Cs\|\}}|\}}}||krdSt|ƒ} |dkr‚t|tƒsDtd| jƒ‚|sPtdƒ‚tdd„|Dƒƒdkrptd|ƒ‚ttt|ƒƒƒ} nt|t ƒs t|tƒr td |ƒ‚| |j |j |  ¡ƒ} |d ksÈ|d krÐ| |kS|d krà| |kS|d krð| |kS|dk�r| |kS|dk�r| |kS|dk�r&| |kS|dk�r8| |kS|dk�rJ| |kStd|dƒ‚dS)NT>rDrCz'%s' object is not a collectionz+Cannot use empty collection as filter valuecSsh|] }t|ƒ’qSr)r�)r-ÚitemrrrÚ ®sz=ParquetPartitions.filter_accepts_partition..r_z8All elements of the collection '%s' must be of same typez-Op '%s' not supported with a collection valuer<r=r>r?r@rArBrCrDz+'%s' is not a valid operator in predicates.) r�r$rrðr„r2r+ÚnextÚiterr'rÙrÚZas_py) r[Úpart_keyÚfilterrôZp_columnZ p_value_indexZf_columnr:Zf_valueZf_typeZp_valuerrrÚfilter_accepts_partition sZ  ÿÿÿ ÿ      ÿz*ParquetPartitions.filter_accepts_partitionN) r„r…r†r^rîrïr»rÎrãrûrrrrrëgs rëc@s>eZdZddd„Zdd„Zd d „Zd d „Zd d„Zdd„ZdS)ÚParquetManifestNú/Úhiver_cCs t||ƒ\}}||_||_||_t|ƒ|_||_tƒ|_g|_ ||_ t j |d�|_ d|_d|_| d|jg¡|j jdd„d�|jdkr’|j|_|j  ¡dS)N)Ú max_workersrcSs|jSr)r©ÚpiecerrrÚæóz*ParquetManifest.__init__..)rf)r"rrÈÚpathseprÚdirpathÚpartition_schemerërÛÚpiecesÚ_metadata_nthreadsrZThreadPoolExecutorÚ _thread_poolÚcommon_metadata_pathÚ metadata_pathÚ _visit_levelÚsortÚshutdown)r[rrÈrrrÚmetadata_nthreadsrrrr^Ñs& ÿ zParquetManifest.__init__c sìˆj}t| ˆ¡ƒ\}}}g}|D]P} ˆj ˆ| f¡} |  d¡rH| ˆ_q"|  d¡rZ| ˆ_q"ˆ | ¡rhq"q"|  | ¡q"‡‡fdd„|Dƒ} |  ¡|   ¡t |ƒdkrÀt | ƒdkrÀt d  ˆ¡ƒ‚n(t | ƒdkr܈ || |¡n ˆ ||¡dS)NZ_common_metadataÚ _metadatacs$g|]}t|ƒsˆj ˆ|f¡‘qSr)Ú_is_private_directoryrrbrå©Ú base_pathr[rrrJsþz0ParquetManifest._visit_level..rz,Found files in an intermediate directory: {})rr÷ÚwalkrrbÚendswithr r Ú_should_silently_excluderMr r+r2rGÚ_visit_directoriesÚ _push_pieces) r[rôrÚ part_keysÚfsr‰Ú directoriesÚfilesZfiltered_filesrÚ full_pathZfiltered_directoriesrrrr îs0     ÿÿ zParquetManifest._visit_levelcCs0| d¡p.| d¡p.| d¡p.| d¡p.|tkS)Nz.crcz _$folder$r`r‰)rÚ startswithÚEXCLUDED_PARQUET_PATHS)r[Ú file_namerrrrs ÿþýüz(ParquetManifest._should_silently_excludec Csšg}|D]~}t||jƒ\}}t|ƒ\}} |j ||| ¡} ||| fg} ||jkrt|j |j|d|| ¡} |  | ¡q| |d|| ¡q|r–t   |¡dS©Nr_) Ú _path_splitrÚ_parse_hive_partitionrÛrãrr Zsubmitr rMrÚwait) r[rôrrZ futures_listrÚheadÚtailr‚rfrÒZ dir_part_keysÚfuturerrrrs    ý z"ParquetManifest._visit_directoriescCs&|jdkrt|ƒStd |j¡ƒ‚dS)Nrþzpartition schema: {})rr#r5rG)r[ÚdirnamerrrÚ_parse_partition+s  ÿz ParquetManifest._parse_partitioncs ˆj ‡‡fdd„|Dƒ¡dS)Ncsg|]}t|ˆˆjd�‘qS))rÊrÈ)rÄrÈ©r-r©rr[rrrJ3sþÿz0ParquetManifest._push_pieces..)rr)r[rrrr+rr2sýzParquetManifest._push_pieces)NNrýrþr_) r„r…r†r^r rrr)rrrrrrüÏsÿ !rücCs"d|krtd |¡ƒ‚| dd¡S)Nr<z3Directory name did not appear to be a partition: {}r_)r2rGÚsplit)Úvaluerrrr#:s ÿr#cCs,tj |¡\}}| d¡s$| d¡o*d|kS)Nr‰r`r<)Úosrr,r)rær‰r&rrrrAsrcCs:| |¡d}|d|…||d…}}| |¡}||fSr!)ÚrfindÚrstrip)rÚseprer%r&rrrr"Fs r"Z_SUCCESSc@seZdZdZdS)Ú_ParquetDatasetMetadata)rr]rTrXrSN)r„r…r†Ú __slots__rrrrr2Psr2cCsD|jdk r(t|jtjƒs(|jj|dd�}t|||j|j|j|j d�S)NrÅrÆ)rUr]rTrXrS) rr$r rrWrRr]rTrXrS)rLrÚmetarrrÚ_open_dataset_fileUs  ÿúr5a€read_dictionary : list, default None List of names or column paths (for nested types) to read directly as DictionaryArray. Only supported for BYTE_ARRAY storage. To read a flat column as dictionary-encoded pass the column name. For nested types, you must pass the full column "path", which could be something like level1.level2.list.item. Refer to the Parquet file's schema to obtain the paths. memory_map : bool, default False If the source is a file path, use a memory map to read file, which can improve performance in some environments. buffer_size : int, default 0 If positive, perform read buffering when deserializing individual column chunks. Otherwise IO calls are unbuffered. partitioning : Partitioning or str or list of str, default "hive" The partitioning scheme for a partitioned dataset. The default of "hive" assumes directory names with key=value pairs like "/year=2009/month=11". In addition, a scheme like "/2009/11" is also supported, in which case you need to specify the field names or a full schema. See the ``pyarrow.dataset.partitioning()`` function for more details.c @s¬eZdZd ee¡Zddd „Zd d d „Zd d „Z dd„Z dd„Z d!dd„Z dd„Z dd„Zdd„Zee d¡ƒZee d¡ƒZee d¡ƒZee d¡ƒZee d¡ƒZdS)"ÚParquetDatasetaP Encapsulates details of reading a complete Parquet dataset possibly consisting of multiple files and partitions in subdirectories. Parameters ---------- path_or_paths : str or List[str] A directory name, single file name, or list of file names. filesystem : FileSystem, default None If nothing passed, paths assumed to be found in the local on-disk filesystem. metadata : pyarrow.parquet.FileMetaData Use metadata obtained elsewhere to validate file schemas. schema : pyarrow.parquet.Schema Use schema obtained elsewhere to validate file schemas. Alternative to metadata parameter. split_row_groups : bool, default False Divide files into pieces for each row group in the file. validate_schema : bool, default True Check that individual file schemas are all the same / compatible. filters : List[Tuple] or List[List[Tuple]] or None (default) Rows which do not match the filter predicate will be removed from scanned data. Partition keys embedded in a nested directory structure will be exploited to avoid loading files at all if they contain no matching rows. If `use_legacy_dataset` is True, filters can only reference partition keys and only a hive-style directory structure is supported. When setting `use_legacy_dataset` to False, also within-file level filtering and different partitioning schemes are supported. {1} metadata_nthreads: int, default 1 How many threads to allow the thread pool which is used to read the dataset metadata. Increasing this is helpful to read partitioned datasets. {0} use_legacy_dataset : bool, default True Set to False to enable the new code path (experimental, using the new Arrow Dataset API). Among other things, this allows to pass `filters` for all columns and not only the partition keys, enables different partitioning schemes, etc. NFTr_rrþcCsN| dkrt|tƒrd} nd} | s@t|||| | | | |||||d� St |¡}|S)NFT) rr6Ú partitioningrTr]rSrirUÚsplit_row_groupsÚvalidate_schemar)r$rÚ_ParquetDatasetV2ÚobjectÚ__new__)ÚclsÚ path_or_pathsrrirUr8r9r6rrTr]rSr7Úuse_legacy_datasetr[rrrr<¥s& ö zParquetDataset.__new__c Cst| dkrtdƒ‚tƒ|_|}t|tƒr.|d}t||ƒ\|j_}t|tƒr\dd„|Dƒ|_n t|ƒ|_| |j_ | |j_ | |j_ t ||j|t t|jƒd�\|_|_|_|_|jdk rÞ|j |j¡�}t|| d�|j_W5QRXnd|j_|dk�r&|jdk �r&|j |j¡�}t|| d�|_W5QRXn||_||_||_|�rFtdƒ‚|dk �rbt|ƒ}| |¡|�rp| ¡dS) NrþzVOnly "hive" for hive-like partitioning is supported when using use_legacy_dataset=TruercSsg|] }t|ƒ‘qSr)rr*rrrrJÑsz+ParquetDataset.__init__..)rrÈ©r]z$split_row_groups not yet implemented)r2r2rr$r3r"rrcrrTr]rSÚ_make_manifestrr5rrÛr r rWÚ read_metadatarXrUrir8r5r;Ú_filterÚvalidate_schemas)r[r>rrirUr8r9r6rrTr]rSr7r?Za_pathr‰r.rrrr^ÁsZÿ    þý þ  zParquetDataset.__init__cCsNt|tƒstdƒ‚|jj|jjkr&dSdD]}t||ƒt||ƒkr*dSq*dS)Nz-`other` must be an instance of ParquetDatasetF) rcr]rrÛr r rXrUrirSr8T)r$r6rðrÚ __class__r¯)r[rÍÚproprrrr»þs zParquetDataset.equalscCs*z | |¡WStk r$tYSXdSrrñrÌrrrrÎ s zParquetDataset.__eq__cCsØ|jdkr>|jdkr>|jdk r*|jj|_qR|jd ¡j|_n|jdkrR|jj|_|j ¡}|jdk r–|jjD]&}| |¡dkrn| |¡}|  |¡}qn|jD]6}| ¡}|j ¡}|j |dd�sœt d  |||¡ƒ‚qœdS)NréÿÿÿÿFr·z-Schema in {!s} was different. {!s} vs {!s}) rUrirXrrÔÚto_arrow_schemarÛríÚget_field_indexÚremover»r2rG)r[Zdataset_schemaZpartition_nameZ field_idxrZ file_metadataZ file_schemarrrrDs*           ýzParquetDataset.validate_schemasc Csng}|jD]"}|j|||j|d�}| |¡q t |¡}|rj| ¡}|jjpNi} |rjd| krj|  d|i¡}|S)aì Read multiple Parquet files as a single pyarrow.Table. Parameters ---------- columns : List[str] Names of columns to read from the file. use_threads : bool, default True Perform multi-threaded column reads use_pandas_metadata : bool, default False Passed through to each dataset piece. Returns ------- pyarrow.Table Content of the file as a table (of columns). )rsrprÛrmr|) rrzrÛrMrçZ concat_tablesÚ_get_common_pandas_metadatarirUÚreplace_schema_metadata) r[rsrprmZtablesrr•Zall_datarXZcurrent_metadatarrrrz/s" þ    ÿzParquetDataset.readcKs|jfddi|—ŽS)a Read dataset including pandas metadata, if any. Other arguments passed through to ParquetDataset.read, see docstring for further details. Returns ------- pyarrow.Table Content of the file as a table (of columns). rmT©rz©r[rµrrrÚ read_pandasWs zParquetDataset.read_pandascCs"|jdkrdS|jj}| dd¡S)Nr|)rXrUÚget)r[rÃrrrrKcs z*ParquetDataset._get_common_pandas_metadatacs<|jj‰‡fdd„‰‡‡fdd„‰‡fdd„|jDƒ|_dS)Ncst‡‡fdd„t|jƒDƒƒS)Nc3s|]\}}ˆ|ˆ|ƒVqdSrr)r-rôrù)Úaccepts_filterrúrrr/nsÿzEParquetDataset._filter..one_filter_accepts..)r4rarÊ)rrú)rQ)rúrÚone_filter_acceptsmsÿz2ParquetDataset._filter..one_filter_acceptscst‡‡fdd„ˆDƒƒS)Nc3s&|]}t‡‡fdd„|DƒƒVqdS)c3s|]}ˆˆ|ƒVqdSrrr,©rRrrrr/rszOParquetDataset._filter..all_filters_accept...N)r4)r-r8rSrrr/rsÿzEParquetDataset._filter..all_filters_accept..)r1r)r6rRrrÚall_filters_acceptqsÿz2ParquetDataset._filter..all_filters_acceptcsg|]}ˆ|ƒr|‘qSrr)r-Úp)rTrrrJusz*ParquetDataset._filter..)rÛrûr)r[r6r)rQrTr6rRrrCjs zParquetDataset._filterz _metadata.fsz_metadata.memory_mapz_metadata.read_dictionaryz_metadata.common_metadataz_metadata.buffer_size) NNNNFTNr_NFrrþN) NNNFTNr_NFrrþT)NTF)r„r…r†rGÚ_read_docstring_commonÚ_DNF_filter_docr‡r<r^r»rÎrDrzrOrKrCrˆrNÚ attrgetterrr]rTrXrSrrrrr6ysX(Ø*ü ü = (  ÿÿr6rýr_c CsÜd}d}d}t|tƒr*t|ƒdkr*|d}t|ƒrp| |¡rpt|||t|ddƒ|d�}|j}|j}|j } |j }n`t|tƒs€|g}t|ƒdkr”t dƒ‚g} |D]2} |  | ¡s¸t d | ¡ƒ‚t| |d�} |  | ¡qœ| |||fS) Nr_rrrý)rrÈrrz Must pass at least one file pathzPassed non-file path: {})rÈ)r$r3r+rÚisdirrür¯r r rrÛr2ÚisfileÚOSErrorrGrÄrM) r>rrrrÈrÛr r ÚmanifestrrrrrrrA‚s8 ý   ÿ  rAc@sDeZdZdZddd„Zedd„ƒZdd d „Zd d „Zedd„ƒZ dS)r:zC ParquetDataset shim using the Dataset API under the hood. NrþFc KsÈddlm} dD]*\} } | | kr| | | k rtd | ¡ƒ‚qi} |rR| jd|d�|dk rf| j|d�||_|ovt|ƒ|_|dk r�t||d�}n|dkr¦|r¦t |d�}d}t |t ƒrÊt |ƒdkrÈ|d}nht |ƒ�r.t|ƒ}|dk�rzt |¡\}}Wn tk �rt |d�}YnX| |¡j�r2|}n|}|dk �r„d|_| jdd �| j| d �}| ||¡}| j|g|j||jd �|_dSd |_| j| d �}|d k�r®| jjdd�}| j|||||d�|_dS)Nr))riN)rUN)r8F)r9T)rr_z;Keyword '{0}' is not yet supported with the new Dataset APIT)Zuse_buffered_streamrS)Zdictionary_columns)Zuse_mmapr_)Z!enable_parallel_column_conversion)Ú read_options)rirGrFrþ)Zinfer_dictionary)rrGr7Úignore_prefixes)rKrLr2rGÚupdateÚ_filtersrQÚ_filter_expressionrrr$r3r+rr'rZfrom_uriZ get_file_infoÚis_fileÚ"_enable_parallel_column_conversionÚParquetFileFormatZ make_fragmentZFileSystemDatasetZphysical_schemarÚ_datasetZHivePartitioningÚdiscover)r[r>rr6r7rTrSr]r^rµrHÚkeywordÚdefaultr]Z single_fileÚparquet_formatÚfragmentrrrr^­s~  ÿÿÿ ÿ       ÿ     ý  ÿýz_ParquetDatasetV2.__init__cCs|jjSr)rerirhrrrrisz_ParquetDatasetV2.schemaTcCs¨|jj}|dk rJ|rJ|rJd|krJdd„t|ƒDƒ}|tt|ƒt|ƒƒ}|jrX|rXd}|jj||j|d�}|r¤|r¤d|kr¤|jjp†i}|  d|di¡|  |¡}|S)a¾ Read (multiple) Parquet files as a single pyarrow.Table. Parameters ---------- columns : List[str] Names of columns to read from the dataset. The partition fields are not automatically included (in contrast to when setting ``use_legacy_dataset=True``). use_threads : bool, default True Perform multi-threaded column reads. use_pandas_metadata : bool, default False If True and file has custom pandas schema metadata, ensure that index columns are also loaded. Returns ------- pyarrow.Table Content of the file as a table (of columns). Nr|cSsg|]}t|tƒs|‘qSr)r$r}©r-r9rrrrJ s ÿz*_ParquetDatasetV2.read..F)rsrúrp) rirUr€r3rìrcreZto_tablerar_rL)r[rsrprmrUrƒr•Z new_metadatarrrrzs*  ÿþ   z_ParquetDatasetV2.readcKs|jfddi|—ŽS)z£ Read dataset including pandas metadata, if any. Other arguments passed through to ParquetDataset.read, see docstring for further details. rmTrMrNrrrrO;sz_ParquetDatasetV2.read_pandascCst|j ¡ƒSr)r3reZ get_fragmentsrhrrrrBsz_ParquetDatasetV2.pieces)NNrþNNFN)NTF) r„r…r†r‡r^rˆrirzrOrrrrrr:¨sþ T  6r:aœ {0} Parameters ---------- source: str, pyarrow.NativeFile, or file-like object If a string passed, can be a single file name or directory name. For file-like objects, only read a single file. Use pyarrow.BufferReader to read a file contained in a bytes or buffer-like object. columns: list If not None, only these columns will be read from the file. A column name may be a prefix of a nested field, e.g. 'a' will select 'a.b', 'a.c', and 'a.d.e'. use_threads : bool, default True Perform multi-threaded column reads. metadata : FileMetaData If separately computed {1} use_legacy_dataset : bool, default False By default, `read_table` uses the new Arrow Datasets API since pyarrow 1.0.0. Among other things, this allows to pass `filters` for all columns and not only the partition keys, enables different partitioning schemes, etc. Set to True to use the legacy behaviour. ignore_prefixes : list, optional Files matching any of these prefixes will be ignored by the discovery process if use_legacy_dataset=False. This is matched to the basename of a path. By default this is ['.', '_']. Note that discovery happens only if a directory is passed as source. filesystem : FileSystem, default None If nothing passed, paths assumed to be found in the local on-disk filesystem. filters : List[Tuple] or List[List[Tuple]] or None (default) Rows which do not match the filter predicate will be removed from scanned data. Partition keys embedded in a nested directory structure will be exploited to avoid loading files at all if they contain no matching rows. If `use_legacy_dataset` is True, filters can only reference partition keys and only a hive-style directory structure is supported. When setting `use_legacy_dataset` to False, also within-file level filtering and different partitioning schemes are supported. {3} Returns ------- {2} Frþc  Csü| s¨|dk rtdƒ‚zt||| ||| || d�} Wnftk r–|dk rPtdƒ‚| dkr`tdƒ‚t||ƒ\}}|dk r€| |¡}t||||| d�} YnX| j|||d�S| dk r¸tdƒ‚t|ƒrÚt||||| ||| d �}nt||||| d�}|j|||d�S) Nz¥The 'metadata' keyword is no longer supported with the new datasets-based implementation. Specify 'use_legacy_dataset=True' to temporarily recover the old behaviour.)rr7r]rTrSr6r^zWthe 'filters' keyword is not supported when the pyarrow.dataset module is not availablerþz\the 'partitioning' keyword is not supported when the pyarrow.dataset module is not available)rUrTr]rSrÕzMThe 'ignore_prefixes' keyword is only supported when use_legacy_dataset=False)rUr]rTrSrr6r7) r2r:Ú ImportErrorrZopen_input_filerRrzrr6)r\rsrprUrmr]rTrr6rSr7r?r^rLrÚpfrrrÚ read_tablezstÿø ÿÿ þ ÿÿüýÿrnz}Read a Table from Parquet format Note: starting with pyarrow 1.0, the default for `use_legacy_dataset` is switched to False.Ú z�use_pandas_metadata : bool, default False If True and file has custom pandas schema metadata, ensure that index columns are also loadedz=pyarrow.Table Content of the file as a table (of columns)c Cst|||||||d||d� S)NT) rsrprUr6r]rSrmr?r^)rn) r\rsrpr]rUr6rSr?r^rrrrOÌsörOzcRead a Table from Parquet format, also reading DataFrame index values if known in the file metadatazgpyarrow.Table Content of the file as a Table of Columns, including DataFrame indexes as columnsršr›cKs¬| d|¡}|}zNt||jf| || |||| | ||| ||dœ |—Ž�}|j||d�W5QRXWnHtk r¦t|ƒr zt t|ƒ¡Wntj k ržYnX‚YnXdS)NÚ chunk_size) rrŸr�r r¡Úcoerce_timestampsÚdata_page_sizeÚallow_truncated_timestampsr�r¢r£r¤r¦r¸) r©r™rir¼Ú Exceptionrr.rJrÚerror)r•r§r¹rŸr r�r¡r¢rqrsrrr�rr£r¤r¦rµZ use_int96r¬rrrr¼çs@ ÿòñr¼z® Write a Table to Parquet format. Parameters ---------- table : pyarrow.Table where: string or pyarrow.NativeFile row_group_size: int The number of rows per rowgroup {} cCsH| ¡rD| |¡sDz| |¡Wn"tk rB| |¡s>t‚YnXdSr)Z _isfilestoreÚexistsÚmkdirr[rº)rrrrrÚ_mkdir_if_not_existss rxc sØ|sÈddlm}| dd¡}| dd¡} d} | dd¡} | dk rNt|  d¡ƒ‚|dk rdt|  d¡ƒ‚| ¡} | jf|Ž} |dk rˆt|ƒ}d}|rª| |¡j }|j |d d �}|j |||| | ||| d �dSt   ||¡\}}t||ƒ| dd¡} |dk �rft|ƒdk�rf| ¡‰‡fd d „|Dƒ}ˆj|dd�}ˆj |¡}t|ƒdk�rPtdƒ‚|j }|j jD] }||k�r^| | |¡¡}�q^| |¡D]Ø\}}t|tƒ�s¤|f}d dd „t||ƒDƒ¡}tjj||dd�}t|d ||g¡ƒ|�rô||ƒ}n tƒd}d ||g¡}d ||g¡}| |d¡�}t ||fd| i|—ŽW5QRX| dk �rŠ| d !|¡�qŠnn|�rv|dƒ}n tƒd}d ||g¡}| |d¡�}t ||fd| i|—ŽW5QRX| dk �rÔ| d !|¡dS)aïWrapper around parquet.write_table for writing a Table to Parquet format by partitions. For each combination of partition columns and values, a subdirectories are created in the following manner: root_dir/ group1=value1 group2=value1 .parquet group2=value2 .parquet group1=valueN group2=value1 .parquet group2=valueN .parquet Parameters ---------- table : pyarrow.Table root_path : str, pathlib.Path The root directory of the dataset filesystem : FileSystem, default None If nothing passed, paths assumed to be found in the local on-disk filesystem partition_cols : list, Column names by which to partition the dataset Columns are partitioned in the order they are given partition_filename_cb : callable, A callback function that takes the partition key(s) as an argument and allow you to override the partition filename. If nothing is passed, the filename will consist of a uuid. use_legacy_dataset : bool, default True Set to False to enable the new code path (experimental, using the new Arrow Dataset API). This is more efficient when using partition columns, but does not (yet) support `partition_filename_cb` and `metadata_collector` keywords. **kwargs : dict, Additional kwargs for write_table function. See docstring for `write_table` or `ParquetWriter` for more information. Using `metadata_collector` in kwargs allows one to collect the file metadata instances of dataset pieces. The file paths in the ColumnChunkMetaData will be set relative to `root_path`. rNrirpTzGThe '{}' argument is not supported with the new dataset implementation.ržÚpartition_filename_cbrþ)r�)rrGrËrir7rpcsg|] }ˆ|‘qSrrrk©ZdfrrrJ‚sz$write_to_dataset..rs)Zaxisz.No data left to save outside partition columnsrýcSsg|]\}}dj||d�‘qS)z{colname}={value})Zcolnamer-rÑ)r-r‚r(rrrrJ”sÿF)riÚsafez.parquetrœrG)"rKrLr©r2rGrdZmake_write_optionsrÚselectrir7Z write_datasetr r!rxr+Z to_pandasZdroprsÚnamesrJrIÚgroupbyr$ÚtuplerbÚziprŽr–Z from_pandasrrWr¼Z set_file_path)r•Ú root_pathZpartition_colsryrr?rµrHrirpr½ržriZ write_optionsr7Z part_schemarrÊZdata_dfZ data_colsZ subschemar9ràZsubgroupÚsubdirZsubtableÚoutfileÚ relative_pathrr.rrzrÚwrite_to_dataset%sš0   ÿ   ý      ÿÿ ÿ  ÿ   ÿ r…cKsHt||f|Ž}| ¡|dk rDt|ƒ}|D]}| |¡q*| |¡dS)a Write metadata-only Parquet file from schema. This can be used with `write_to_dataset` to generate `_common_metadata` and `_metadata` sidecar files. Parameters ---------- schema : pyarrow.Schema where: string or pyarrow.NativeFile metadata_collector: **kwargs : dict, Additional kwargs for ParquetWriter class. See docstring for `ParquetWriter` for more information. Examples -------- Write a dataset and collect metadata information. >>> metadata_collector = [] >>> write_to_dataset( ... table, root_path, ... metadata_collector=metadata_collector, **writer_kwargs) Write the `_common_metadata` parquet file without row groups statistics. >>> write_metadata( ... table.schema, root_path / '_common_metadata', **writer_kwargs) Write the `_metadata` parquet file with row groups statistics. >>> write_metadata( ... table.schema, root_path / '_metadata', ... metadata_collector=metadata_collector, **writer_kwargs) N)r™r°rBZappend_row_groupsZwrite_metadata_file)rir§ržrµr¬rUÚmrrrÚwrite_metadata±s$ r‡cCst||d�jS)a# Read FileMetadata from footer of a single Parquet file. Parameters ---------- where : str (filepath) or file-like object memory_map : bool, default False Create memory map when the source is a file path. Returns ------- metadata : FileMetadata r@)rRrU©r§r]rrrrBásrBcCst||d�j ¡S)a# Read effective Arrow schema from Parquet file metadata. Parameters ---------- where : str (filepath) or file-like object memory_map : bool, default False Create memory map when the source is a file path. Returns ------- schema : pyarrow.Schema r@)rRrirHrˆrrrÚ read_schemaòsr‰)T)N)rýr_N) NTNFFNNNrrþFN)NTFNNrTN)NršTr›TNNFNNNNFrš)NNNT)N)F)F)SÚ collectionsrZ concurrentrÚ functoolsrrrÀÚcollections.abcrÚnumpyr×r.ÚrerNÚ urllib.parserZpyarrowrŽZ pyarrow.librçZpyarrow._parquetr«rrr r r r r Z pyarrow.fsrrrrrr Z pyarrow.utilrrrrrr"r*r;rWrQrRÚcompilerŠrŒr“r˜r¾r™r€rÄrÝrërür#rr"rr2r5rVr6rAr:Z_read_table_docstringrnrGrbr‡rOr¼rxr…r‡rBr‰rrrrÚsà     $    6z  Bt  ;hk  ÿ &!2ü Dÿõþ ù ö ) ö þ  0