So I've been going through the deserialization log...
# daft-dev
r
So I've been going through the deserialization logic (specifically the index deserialization logic) in the
arrow2
crate. There are two options that I'm thinking of in order to allow for larger strings inside of our Utf8Array: 1. Modify the original
arrow2::Utf8Array
type to not be parameterized over possible index types. Currently,
Utf8Array
is defined as
struct UtfArray<O: Offset>
where
impl Offset for i32 { .. }
and
impl Offset for i64 { .. }
exist. I could change it to simply always have an internal index of `i64`s and thus remove any parameterization over
Utf8Array
. This will require modifications to the
Utf8Array
type inside of the
arrow2
crate, which I don't think will be ideal. 2. The other option is that we always convert all utf8 arrays inside of
daft-core
to be of datatype
arrow2::DataType::LargeUtf8
(which uses
i64
offsets internally) instead of being
arrow2::DataType::Utf8
. This will then always follow the
i64
deserialization logic inside of
arrow2
. I think the 2nd option will require the least number of changes and will also not require me to modify the
arrow2
crate, which we've pulled into tree. Any thoughts?
j
(2) sounds good yeah
☝️ 1
❤️ 1
s
(2) sounds better. (1) is probably not possible since we still need
i32
UtfArrays since we might be handled some over from the python side via FFI.
❤️ 1
r
Interesting. I can't find any instances of
Utf8Array<i32>
inside of
daft-core
at all... I only see
Utf8Array<i64>
.
s
We do the conversion here The type gets converted here
r
We're running in an
ArrowCapacityError
after deserializing directly into an i64-based offset buffer @Sammy Sidhu.
This exception is being sourced from
pyarrow
.
This seems like an error originating from storage, since PyArrow doesn't support this length of strings when the data-at-rest is deserialized into memory (even if the offsets are i64 encoded) @Sammy Sidhu.
Seems like this is just a natural difference between the two different encoding paradigms of Parquet and Arrow.
I don't think PyArrow does any chunking internally. This may have to be manually done...
Repro steps: 1. checkout my branch:
bug/parquet-large-files-reader
2. make build-release 3.
daft.read_parquet('<hf://datasets/AlgorithmicResearchGroup/arxiv_research_code/').write('hf>')
The stack trace error that I'm getting I will send in a follow-up screenshot.
j
We’re running in an
ArrowCapacityError
after deserializing directly into an i64-based offset buffer
The image indicates that it’s an i32 offsets buffer I think, given that that’s i32::max
Seems at some point the pyarrow writer is using small string types. You should check if 1. the data is getting passed in correctly as large strings 2. PyArrow isn’t doing something dumb under the hood, or maybe we need to hint it or something to use large strings?
r
I have a small fix in there to immediately deserialize into an i64 offset buffer.
Copy code
---------------------------------------------------------------------------
ArrowCapacityError                        Traceback (most recent call last)
Cell In[2], line 1
----> 1 daft.read_parquet('<hf://datasets/AlgorithmicResearchGroup/arxiv_research_code/').write_parquet('hf3>')

File ~/home/eventual-repos/Daft/daft/api_annotations.py:26, in DataframePublicAPI.<locals>._wrap(*args, **kwargs)
     24 type_check_function(func, *args, **kwargs)
     25 timed_method = time_df_method(func)
---> 26 return timed_method(*args, **kwargs)

File ~/home/eventual-repos/Daft/daft/analytics.py:198, in time_df_method.<locals>.tracked_method(*args, **kwargs)
    195 @functools.wraps(method)
    196 def tracked_method(*args, **kwargs):
    197     if _ANALYTICS_CLIENT is None:
--> 198         return method(*args, **kwargs)
    200     start = time.time()
    201     try:

File ~/home/eventual-repos/Daft/daft/dataframe/dataframe.py:554, in DataFrame.write_parquet(self, root_dir, compression, partition_cols, io_config)
    552 # Block and write, then retrieve data
    553 write_df = DataFrame(builder)
--> 554 write_df.collect()
    555 assert write_df._result is not None
    557 if len(write_df) > 0:
    558     # Populate and return a new disconnected DataFrame

File ~/home/eventual-repos/Daft/daft/api_annotations.py:26, in DataframePublicAPI.<locals>._wrap(*args, **kwargs)
     24 type_check_function(func, *args, **kwargs)
     25 timed_method = time_df_method(func)
---> 26 return timed_method(*args, **kwargs)

File ~/home/eventual-repos/Daft/daft/analytics.py:198, in time_df_method.<locals>.tracked_method(*args, **kwargs)
    195 @functools.wraps(method)
    196 def tracked_method(*args, **kwargs):
    197     if _ANALYTICS_CLIENT is None:
--> 198         return method(*args, **kwargs)
    200     start = time.time()
    201     try:

File ~/home/eventual-repos/Daft/daft/dataframe/dataframe.py:2483, in DataFrame.collect(self, num_preview_rows)
   2470 @DataframePublicAPI
   2471 def collect(self, num_preview_rows: Optional[int] = 8) -> "DataFrame":
   2472     """Executes the entire DataFrame and materializes the results
   2473
   2474     .. NOTE::
   (...)
   2481         DataFrame: DataFrame with materialized results.
   2482     """
-> 2483     self._materialize_results()
   2485     assert self._result is not None
   2486     dataframe_len = len(self._result)

File ~/home/eventual-repos/Daft/daft/dataframe/dataframe.py:2465, in DataFrame._materialize_results(self)
   2463 context = get_context()
   2464 if self._result is None:
-> 2465     self._result_cache = context.runner().run(self._builder)
   2466     result = self._result
   2467     assert result is not None

File ~/home/eventual-repos/Daft/daft/runners/pyrunner.py:404, in PyRunner.run(self, builder)
    403 def run(self, builder: LogicalPlanBuilder) -> PartitionCacheEntry:
--> 404     results = list(self.run_iter(builder))
    406     result_pset = LocalPartitionSet()
    407     for i, result in enumerate(results):

File ~/home/eventual-repos/Daft/daft/runners/pyrunner.py:469, in PyRunner.run_iter(self, builder, results_buffer_size)
    467 with profiler("profile_PyRunner.run_{datetime.now().isoformat()}.json"):
    468     results_gen = self._physical_plan_to_partitions(execution_id, tasks)
--> 469     yield from results_gen

File ~/home/eventual-repos/Daft/daft/runners/pyrunner.py:642, in PyRunner._physical_plan_to_partitions(self, execution_id, plan)
    640 for done_future in done_set:
    641     done_task = local_futures_to_task.pop(done_future)
--> 642     materialized_results = done_future.result()
    644     pbar.mark_task_done(done_task)
    645     del self._inflight_futures[(execution_id, done_task.id())]

File /Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/lib/python3.9/concurrent/futures/_base.py:438, in Future.result(self, timeout)
    436     raise CancelledError()
    437 elif self._state == FINISHED:
--> 438     return self.__get_result()
    440 self._condition.wait(timeout)
    442 if self._state in [CANCELLED, CANCELLED_AND_NOTIFIED]:

File /Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/lib/python3.9/concurrent/futures/_base.py:390, in Future.__get_result(self)
    388 if self._exception:
    389     try:
--> 390         raise self._exception
    391     finally:
    392         # Break a reference cycle with the exception in self._exception
    393         self = None

File /Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/lib/python3.9/concurrent/futures/thread.py:52, in _WorkItem.run(self)
     49     return
     51 try:
---> 52     result = self.fn(*self.args, **self.kwargs)
     53 except BaseException as exc:
     54     self.future.set_exception(exc)

File ~/home/eventual-repos/Daft/daft/runners/pyrunner.py:680, in PyRunner.build_partitions(self, instruction_stack, partitions, final_metadata)
    673 def build_partitions(
    674     self,
    675     instruction_stack: list[Instruction],
    676     partitions: list[MicroPartition],
    677     final_metadata: list[PartialPartitionMetadata],
    678 ) -> list[MaterializedResult[MicroPartition]]:
    679     for instruction in instruction_stack:
--> 680         partitions = instruction.run(partitions)
    682     results: list[MaterializedResult[MicroPartition]] = [
    683         PyMaterializedResult(part, PartitionMetadata.from_table(part).merge_with_partial(partial))
    684         for part, partial in zip(partitions, final_metadata)
    685     ]
    686     return results

File ~/home/eventual-repos/Daft/daft/execution/execution_step.py:359, in WriteFile.run(self, inputs)
    358 def run(self, inputs: list[MicroPartition]) -> list[MicroPartition]:
--> 359     return self._write_file(inputs)

File ~/home/eventual-repos/Daft/daft/execution/execution_step.py:363, in WriteFile._write_file(self, inputs)
    361 def _write_file(self, inputs: list[MicroPartition]) -> list[MicroPartition]:
    362     [input] = inputs
--> 363     partition = self._handle_file_write(
    364         input=input,
    365     )
    366     return [partition]

File ~/home/eventual-repos/Daft/daft/execution/execution_step.py:378, in WriteFile._handle_file_write(self, input)
    377 def _handle_file_write(self, input: MicroPartition) -> MicroPartition:
--> 378     return table_io.write_tabular(
    379         input,
    380         path=self.root_dir,
    381         schema=self.schema,
    382         file_format=self.file_format,
    383         compression=self.compression,
    384         partition_cols=self.partition_cols,
    385         io_config=self.io_config,
    386     )

File ~/home/eventual-repos/Daft/daft/table/table_io.py:512, in write_tabular(table, file_format, path, schema, partition_cols, compression, io_config)
    509     target_row_groups = max(math.ceil(size_bytes / TARGET_ROW_GROUP_SIZE / inflation_factor), 1)
    510     rows_per_row_group = max(min(math.ceil(num_rows / target_row_groups), rows_per_file), 1)
--> 512     _write_tabular_arrow_table(
    513         arrow_table=part_table,
    514         schema=part_table.schema,
    515         full_path=part_path,
    516         format=format,
    517         opts=opts,
    518         fs=fs,
    519         rows_per_file=rows_per_file,
    520         rows_per_row_group=rows_per_row_group,
    521         create_dir=is_local_fs,
    522         file_visitor=visitors.visitor(i),
    523     )
    525 return visitors.to_metadata()

File ~/home/eventual-repos/Daft/daft/table/table_io.py:867, in _write_tabular_arrow_table(arrow_table, schema, full_path, format, opts, fs, rows_per_file, rows_per_row_group, create_dir
    864     ERROR_MSGS = ("InvalidPart", "curlCode: 28, Timeout was reached")
    865     return isinstance(e, OSError) and any(err_str in str(e) for err_str in ERROR_MSGS)
--> 867 _retry_with_backoff(
    868     write_dataset,
    869     full_path,
    870     retry_error=retry_error,
    871 )

File ~/home/eventual-repos/Daft/daft/table/table_io.py:798, in _retry_with_backoff(func, path, retry_error, num_tries, jitter_ms, max_backoff_ms)
    796 for attempt in range(num_tries):
    797     try:
--> 798         return func()
    799     except Exception as e:
    800         if retry_error(e):

File ~/home/eventual-repos/Daft/daft/table/table_io.py:848, in _write_tabular_arrow_table.<locals>.write_dataset()
    847 def write_dataset():
--> 848     pads.write_dataset(
    849         arrow_table,
    850         schema=schema,
    851         base_dir=full_path,
    852         basename_template=basename_template,
    853         format=format,
    854         partitioning=None,
    855         file_options=opts,
    856         file_visitor=file_visitor,
    857         use_threads=True,
    858         existing_data_behavior="overwrite_or_ignore",
    859         filesystem=fs,
    860         **kwargs,
    861     )

File ~/home/eventual-repos/Daft/.venv/lib/python3.9/site-packages/pyarrow/dataset.py:1030, in write_dataset(data, base_dir, basename_template, format, partitioning, partitioning_flavor, schema, filesystem, file_options, use_threads, max_partitions, max_open_files, max_rows_per_file, min_rows_per_group, max_rows_per_group, file_visitor, existing_data_behavior, create_dir)
   1027         raise ValueError("Cannot specify a schema when writing a Scanner")
   1028     scanner = data
-> 1030 _filesystemdataset_write(
   1031     scanner, base_dir, basename_template, filesystem, partitioning,
   1032     file_options, max_partitions, file_visitor, existing_data_behavior,
   1033     max_open_files, max_rows_per_file,
   1034     min_rows_per_group, max_rows_per_group, create_dir
   1035 )

File ~/home/eventual-repos/Daft/.venv/lib/python3.9/site-packages/pyarrow/_dataset.pyx:4010, in pyarrow._dataset._filesystemdataset_write()

File ~/home/eventual-repos/Daft/.venv/lib/python3.9/site-packages/pyarrow/error.pxi:91, in pyarrow.lib.check_status()

ArrowCapacityError: array cannot contain more than 2147483646 bytes, have 2148063122
@jay I don't know this if this is an egress issue, actually. I tested my hypothesis by running this script:
Copy code
import daft

df = None

with open('random_data.txt', encoding='utf-8') as file:
    contents = file.read()
    df = daft.from_pydict({'a': [contents]})

assert df
df = df.collect()
(which notably does not perform any egress [i.e., parquet writes]) and I'm still getting the same
ArrowCapacityError
. The stacktrace is below:
Copy code
(.venv) rabh@Raunaks-MBP ~/h/e/Daft (bug/parquet-large-files-reader)> python noise.py
Traceback (most recent call last):
  File "/Users/rabh/home/eventual-repos/Daft/noise.py", line 7, in <module>
    df = daft.from_pydict({'a': [contents]})
  File "/Users/rabh/home/eventual-repos/Daft/daft/api_annotations.py", line 39, in _wrap
    return timed_func(*args, **kwargs)
  File "/Users/rabh/home/eventual-repos/Daft/daft/analytics.py", line 225, in tracked_fn
    return fn(*args, **kwargs)
  File "/Users/rabh/home/eventual-repos/Daft/daft/convert.py", line 80, in from_pydict
    return DataFrame._from_pydict(data)
  File "/Users/rabh/home/eventual-repos/Daft/daft/dataframe/dataframe.py", line 452, in _from_pydi
ct
    data_micropartition = MicroPartition.from_pydict(data)
  File "/Users/rabh/home/eventual-repos/Daft/daft/table/micropartition.py", line 107, in from_pydi
ct
    table = Table.from_pydict(data)
  File "/Users/rabh/home/eventual-repos/Daft/daft/table/table.py", line 135, in from_pydict
    series = item_to_series(k, v)
  File "/Users/rabh/home/eventual-repos/Daft/daft/series.py", line 636, in item_to_series
    series = Series.from_pylist(item, name)
  File "/Users/rabh/home/eventual-repos/Daft/daft/series.py", line 91, in from_pylist
    arrow_array = pa.array(data)
  File "pyarrow/array.pxi", line 355, in pyarrow.lib.array
  File "pyarrow/array.pxi", line 42, in pyarrow.lib._sequence_to_array
  File "pyarrow/error.pxi", line 154, in pyarrow.lib.pyarrow_internal_check_status
  File "pyarrow/error.pxi", line 91, in pyarrow.lib.check_status
pyarrow.lib.ArrowCapacityError: array cannot contain more than 2147483646 bytes, have 3221225472
(Where
random_data.txt
is a 3gb file of random utf-8 characters).
j
Same error message, very different backtrace. That’s an ingress issue
r
Different backtrace, but does this not point to some similar underlying error?
To be fair, not sure how pyarrow is being used in ingress. I believe that in egress, we use it to write to a parquet file.
j
Same error, different reason 1. Ingress is happening because we use PyArrow to parse the Python string, and it’s defaulting to small string 2. Egress is happening because we use PyArrow to write the data, and it’s doing something weird there Again, I would just breakpoint in the write path and figure out what’s happening there
r
I've been breakpointing in various functions that appeared in the original stack trace. I know where the exception is being thrown. Just trying more to see what the solution here can be.
I placed a breakpoint at the interface between daft and pyarrow (right when the write happens). This intersection happens at
daft/table/table_io.py
inside of the function named
_write_tabular_arrow_table
on line 848. I placed a try-except around here and a breakpoint when the error is thrown.
Copy code
def write_dataset():
        try:
            pads.write_dataset(            # line 848
                arrow_table,
                schema=schema,
                base_dir=full_path,
                basename_template=basename_template,
                format=format,
                partitioning=None,
                file_options=opts,
                file_visitor=file_visitor,
                use_threads=True,
                existing_data_behavior="overwrite_or_ignore",
                filesystem=fs,
                **kwargs,
            )
        except Exception as e:
            breakpoint()
After hitting the breakpoint, I printed out the
arrow_table
and
schema
. This is the printout of
schema
. The
code
column is what is causing the problem in our situation. But all of them are `large_string`'s.
Copy code
(Pdb) schema
repo: large_string
file: large_string
code: large_string
file_length: int64
avg_line_length: double
max_line_length: int64
extension_type: large_string
This is a place that I've been throwing a lot of breakpoints into. It's where each materialized micropartition is materialized and collected. Link: https://github.com/Eventual-Inc/Daft/blob/bug/parquet-large-files-reader/daft/runners/pyrunner.py#L403 Looking into what that generator is yielding currently. The generator that is generating those values is: https://github.com/Eventual-Inc/Daft/blob/bug/parquet-large-files-reader/daft/runners/pyrunner.py#L523
j
(Synced offline and resolved)