Raunak Bhagat
10/20/2024, 11:27 PMarrow2 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?jay
10/20/2024, 11:36 PMSammy Sidhu
10/21/2024, 6:29 AMi32 UtfArrays since we might be handled some over from the python side via FFI.Raunak Bhagat
10/21/2024, 7:32 PMUtf8Array<i32> inside of daft-core at all... I only see Utf8Array<i64>.Raunak Bhagat
10/22/2024, 10:46 PMArrowCapacityError after deserializing directly into an i64-based offset buffer @Sammy Sidhu.Raunak Bhagat
10/22/2024, 10:46 PMpyarrow.Raunak Bhagat
10/22/2024, 10:57 PMRaunak Bhagat
10/22/2024, 11:13 PMRaunak Bhagat
10/22/2024, 11:13 PMRaunak Bhagat
10/22/2024, 11:42 PMRaunak Bhagat
10/22/2024, 11:52 PMbug/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.jay
10/22/2024, 11:54 PMWe’re running in anThe image indicates that it’s an i32 offsets buffer I think, given that that’s i32::maxafter deserializing directly into an i64-based offset bufferArrowCapacityError
jay
10/22/2024, 11:56 PMRaunak Bhagat
10/22/2024, 11:57 PMRaunak Bhagat
10/23/2024, 12:12 AM---------------------------------------------------------------------------
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 2148063122Raunak Bhagat
10/23/2024, 1:15 AMimport 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:
(.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 3221225472Raunak Bhagat
10/23/2024, 1:16 AMrandom_data.txt is a 3gb file of random utf-8 characters).jay
10/23/2024, 1:16 AMRaunak Bhagat
10/23/2024, 1:17 AMRaunak Bhagat
10/23/2024, 1:18 AMjay
10/23/2024, 1:19 AMRaunak Bhagat
10/23/2024, 1:43 AMRaunak Bhagat
10/23/2024, 2:41 AMdaft/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.
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.
(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_stringRaunak Bhagat
10/23/2024, 5:53 AMjay
10/23/2024, 6:57 AM