From 074957cd1f03d5155b28c569ec76a33368e47dd5 Mon Sep 17 00:00:00 2001 From: Yao You Date: Thu, 1 Oct 2026 19:46:10 -0500 Subject: [PATCH 1/6] fix(csv): reject CSV and TSV files that span too many cells Pandas sizes the data-frame by the first record and pads every shorter record to that width. A few-KB file whose first line is a long run of delimiters therefore spans millions of cells; partition_csv() and partition_tsv() render every one of them to HTML and parse it back, at ~400 B per cell. Measure the span by streaming the records with the csv module before pd.read_csv, stopping as soon as it passes CSV_MAX_CELLS (default 5,000,000), and raise UnprocessableEntityError. When the delimiter is left to Pandas (sep=None), sniff it the same way Pandas does, from the first non-blank line, so the measured shape matches what Pandas reads. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 6 +++ test_unstructured/partition/test_csv.py | 53 +++++++++++++++++++++++++ test_unstructured/partition/test_tsv.py | 44 ++++++++++++++++++++ unstructured/__version__.py | 2 +- unstructured/partition/csv.py | 45 +++++++++++++++++++++ unstructured/partition/tsv.py | 6 +++ unstructured/partition/utils/config.py | 5 +++ 7 files changed, 160 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 75720b6c14..f14818b793 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,9 @@ +## 0.27.12 + +### Fixes + +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming its records before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. + ## 0.27.11 ### Fixes diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index b648134d08..4be45dba9c 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -3,7 +3,9 @@ from __future__ import annotations import io +from pathlib import Path +import pandas as pd import pytest from pytest_mock import MockFixture @@ -28,6 +30,7 @@ from unstructured.chunking.title import chunk_by_title from unstructured.cleaners.core import clean_extra_whitespace from unstructured.documents.elements import Table +from unstructured.errors import UnprocessableEntityError from unstructured.partition.csv import _CsvPartitioningContext, partition_csv from unstructured.partition.utils.constants import UNSTRUCTURED_INCLUDE_DEBUG_METADATA @@ -211,6 +214,56 @@ def test_partition_csv_header(): assert table.metadata.text_as_html is not None +# -- cell-count limit ---------------------------------------------------------------------------- + + +def test_partition_csv_rejects_a_wide_first_line_before_pandas_reads_it( + tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) # -- default limit of 5M cells -- + # -- Pandas pads every row out to the first line's 5,000 fields: 25M cells from 15KB -- + file_path = tmp_path / "ragged.csv" + file_path.write_text("h" + "," * 4999 + "\n" + "a\n" * 5000) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="1,001 rows x 5,000 columns"): + partition_csv(str(file_path)) + + read_csv_.assert_not_called() + + +@pytest.mark.parametrize("from_file", [False, True]) +@pytest.mark.parametrize( + ("content", "n_cells"), + [ + # -- the context's restricted sniffer gives up on this one and Pandas sniffs it itself -- + ("h" + "," * 99 + "\n" + "a\n" * 99, 100 * 100), + ("a;b;c\n1;2\n\n4\n", 3 * 3), # -- blank lines are not counted -- + # -- single-column file; Pandas sniffs "a" as the delimiter and reads 4 x 2 cells -- + ("a\nb\nc\nd\n", 4 * 2), + ('"x,\ny",z\n1\n', 2 * 2), # -- quoted delimiter and newline are not counted -- + ], +) +def test_partition_csv_limits_the_cells_the_file_spans( + content: str, n_cells: int, from_file: bool, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + file_path = tmp_path / "table.csv" + file_path.write_text(content) + + def partition(): + if from_file: + with open(file_path, "rb") as f: + return partition_csv(file=f) + return partition_csv(str(file_path)) + + monkeypatch.setenv("CSV_MAX_CELLS", str(n_cells)) + assert len(partition()) == 1 + + monkeypatch.setenv("CSV_MAX_CELLS", str(n_cells - 1)) + with pytest.raises(UnprocessableEntityError, match="CSV_MAX_CELLS"): + partition() + + # ================================================================================================ # UNIT-TESTS # ================================================================================================ diff --git a/test_unstructured/partition/test_tsv.py b/test_unstructured/partition/test_tsv.py index 5276e6006d..8eabe8027b 100644 --- a/test_unstructured/partition/test_tsv.py +++ b/test_unstructured/partition/test_tsv.py @@ -2,6 +2,9 @@ from __future__ import annotations +from pathlib import Path + +import pandas as pd import pytest from pytest_mock import MockFixture @@ -15,6 +18,7 @@ from test_unstructured.unit_utils import assert_round_trips_through_JSON, example_doc_path from unstructured.chunking.title import chunk_by_title from unstructured.documents.elements import Table +from unstructured.errors import UnprocessableEntityError from unstructured.partition.tsv import partition_tsv EXPECTED_FILETYPE = "text/tsv" @@ -159,3 +163,43 @@ def test_partition_tsv_supports_chunking_strategy_while_partitioning(): # The same chunks are returned if chunking elements or chunking during partitioning. assert chunk_elements == chunks + + +# -- cell-count limit ---------------------------------------------------------------------------- + + +@pytest.mark.parametrize("from_file", [False, True]) +def test_partition_tsv_rejects_a_wide_first_line_before_pandas_reads_it( + from_file: bool, tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) # -- default limit of 5M cells -- + # -- Pandas pads every row out to the first line's 5,000 fields: 25M cells from 15KB -- + file_path = tmp_path / "ragged.tsv" + file_path.write_text("h" + "\t" * 4999 + "\n" + "a\n" * 5000) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="1,001 rows x 5,000 columns"): + if from_file: + with open(file_path, "rb") as f: + partition_tsv(file=f) + else: + partition_tsv(str(file_path)) + + read_csv_.assert_not_called() + + +@pytest.mark.parametrize("from_file", [False, True]) +def test_partition_tsv_partitions_a_file_at_the_cell_limit( + from_file: bool, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.setenv("CSV_MAX_CELLS", "6") + file_path = tmp_path / "table.tsv" + file_path.write_text("a\tb\tc\n1\t2\n") + + if from_file: + with open(file_path, "rb") as f: + elements = partition_tsv(file=f) + else: + elements = partition_tsv(str(file_path)) + + assert [e.text for e in elements] == ["a b c 1 2"] diff --git a/unstructured/__version__.py b/unstructured/__version__.py index bef9bb955e..f1527a87a7 100644 --- a/unstructured/__version__.py +++ b/unstructured/__version__.py @@ -1 +1 @@ -__version__ = "0.27.11" # pragma: no cover +__version__ = "0.27.12" # pragma: no cover diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 85c847d948..8d566c51b4 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -1,7 +1,9 @@ from __future__ import annotations +import codecs import contextlib import csv +import itertools from functools import cached_property from typing import IO, Any, Iterator @@ -10,8 +12,10 @@ from unstructured.chunking import add_chunking_strategy from unstructured.common.html_table import HtmlTable from unstructured.documents.elements import Element, ElementMetadata, Table +from unstructured.errors import UnprocessableEntityError from unstructured.file_utils.model import FileType from unstructured.partition.common.metadata import apply_metadata, get_last_modified_date +from unstructured.partition.utils.config import env_config from unstructured.telemetry import partition_runtime_telemetry from unstructured.utils import is_temp_file_path @@ -59,6 +63,8 @@ def partition_csv( ) csv.field_size_limit(CSV_FIELD_LIMIT) + with ctx.open() as file: + check_cell_count(file, ctx.delimiter, ctx.encoding) with ctx.open() as file: read_kw: dict = {"header": ctx.header, "sep": ctx.delimiter, "encoding": ctx.encoding} # sep=None is not supported by the C engine; use Python engine to avoid ParserWarning. @@ -80,6 +86,45 @@ def partition_csv( return [Table(text=html_table.text, metadata=metadata, detection_origin=DETECTION_ORIGIN)] +def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | None) -> None: + """Raise `UnprocessableEntityError` when `file` would span more than `CSV_MAX_CELLS` cells. + + Pandas sizes the data-frame by the first record and pads every shorter record out to that + width, so a tiny file whose first line is a long run of delimiters can span millions of cells. + The span is measured here by streaming the records, before Pandas allocates anything, and the + scan stops as soon as the limit is passed. `file` is read from its current position. + """ + max_cells = env_config.CSV_MAX_CELLS + # -- `errors="replace"` so the scan never fails on a file Pandas would read -- + lines: Iterator[str] = codecs.getreader(encoding or "utf-8")(file, errors="replace") + + if delimiter is None: + # -- With `sep=None` Pandas sniffs the delimiter itself, from the first non-blank line + # -- and without restricting the candidates, so measure with the delimiter it will use. + first_line = next((line for line in lines if line.strip()), "") + try: + delimiter = csv.Sniffer().sniff(first_line).delimiter + except csv.Error: + delimiter = None + lines = itertools.chain([first_line], lines) + + # -- a single-column file has no delimiter; any character absent from each line will do -- + records = csv.reader(lines, delimiter=delimiter or "\n") + + n_rows, n_cols = 0, 0 + for record in records: + if not record: # -- Pandas skips blank lines -- + continue + if n_rows == 0: + n_cols = len(record) + n_rows += 1 + if n_rows * n_cols > max_cells: + raise UnprocessableEntityError( + f"File exceeds the maximum of {max_cells:,} table cells (CSV_MAX_CELLS): it spans" + f" at least {n_rows:,} rows x {n_cols:,} columns." + ) + + class _CsvPartitioningContext: """Encapsulates the partitioning-run details. diff --git a/unstructured/partition/tsv.py b/unstructured/partition/tsv.py index 1a5844eea0..04643772be 100644 --- a/unstructured/partition/tsv.py +++ b/unstructured/partition/tsv.py @@ -13,6 +13,7 @@ spooled_to_bytes_io_if_needed, ) from unstructured.partition.common.metadata import apply_metadata, get_last_modified_date +from unstructured.partition.csv import check_cell_count from unstructured.telemetry import partition_runtime_telemetry DETECTION_ORIGIN: str = "tsv" @@ -44,12 +45,17 @@ def partition_tsv( header = 0 if include_header else None if filename: + with open(filename, "rb") as f: + check_cell_count(f, "\t", None) dataframe = pd.read_csv(filename, sep="\t", header=header) else: assert file is not None # -- Note(scanny): `SpooledTemporaryFile` on Python<3.11 does not implement `.readable()` # -- which triggers an exception on `pd.DataFrame.read_csv()` call. f = spooled_to_bytes_io_if_needed(file) + start = f.tell() + check_cell_count(f, "\t", None) + f.seek(start) dataframe = pd.read_csv(f, sep="\t", header=header) html_table = HtmlTable.from_html_text( diff --git a/unstructured/partition/utils/config.py b/unstructured/partition/utils/config.py index 8e6e15e0a4..191810ea29 100644 --- a/unstructured/partition/utils/config.py +++ b/unstructured/partition/utils/config.py @@ -328,5 +328,10 @@ def PDF_RENDER_MAX_PIXELS_PER_PAGE(self) -> int: """Maximum rendered pixels allowed for a single PDF page""" return self._get_int("PDF_RENDER_MAX_PIXELS_PER_PAGE", 1_000_000_000) + @property + def CSV_MAX_CELLS(self) -> int: + """Maximum `rows x columns` cells a CSV or TSV file may span""" + return self._get_int("CSV_MAX_CELLS", 5_000_000) + env_config = ENVConfig() From 295536d6aa60eb396fe252d1510d4de5848b4d7f Mon Sep 17 00:00:00 2001 From: Yao You Date: Thu, 1 Oct 2026 20:51:48 -0500 Subject: [PATCH 2/6] fix(csv): measure the span without the csv module; normalize line endings Review follow-ups to the CSV/TSV cell limit: - Scan in fixed-size chunks with a quote-aware counter for the first record instead of csv.reader, so a delimiter-heavy record is never built in memory, the csv module's 128 KiB field limit no longer rejects large TSV fields, and no "\n" delimiter is needed (Python 3.13 rejects it). - Count only "\r"/"\n" as line endings. codecs line splitting also broke on "\x0c" and other characters Pandas reads as content, which let a wide first record be measured as one column. - Treat whitespace-only lines as blank the way each Pandas engine does, so a blank first line no longer hides the wide record after it. - Normalize line endings to "\n" before Pandas reads the file. Pandas 2.x's C tokenizer reads a "\r" line ending followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows inside pd.read_csv(). - When the delimiter is sniffed, sniff it once from the normalized first line and pass it to Pandas, so the scan and the read agree; a file with no usable delimiter is read as one column. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- test_unstructured/partition/test_csv.py | 68 +++++-- test_unstructured/partition/test_tsv.py | 12 +- unstructured/partition/csv.py | 230 ++++++++++++++++++++---- unstructured/partition/tsv.py | 9 +- 5 files changed, 269 insertions(+), 52 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index f14818b793..a736a74683 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,7 @@ ### Fixes -- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming its records before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. Line endings are also normalized to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read a `"\r"` line ending followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. When the delimiter is sniffed, it is now sniffed once and passed to Pandas, and a file with no usable delimiter is read as one column. ## 0.27.11 diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index 4be45dba9c..140708931b 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -31,7 +31,7 @@ from unstructured.cleaners.core import clean_extra_whitespace from unstructured.documents.elements import Table from unstructured.errors import UnprocessableEntityError -from unstructured.partition.csv import _CsvPartitioningContext, partition_csv +from unstructured.partition.csv import _CsvPartitioningContext, check_cell_count, partition_csv from unstructured.partition.utils.constants import UNSTRUCTURED_INCLUDE_DEBUG_METADATA EXPECTED_FILETYPE = "text/csv" @@ -226,7 +226,7 @@ def test_partition_csv_rejects_a_wide_first_line_before_pandas_reads_it( file_path.write_text("h" + "," * 4999 + "\n" + "a\n" * 5000) read_csv_ = mocker.patch.object(pd, "read_csv") - with pytest.raises(UnprocessableEntityError, match="1,001 rows x 5,000 columns"): + with pytest.raises(UnprocessableEntityError, match="rows x 5,000 columns"): partition_csv(str(file_path)) read_csv_.assert_not_called() @@ -234,21 +234,29 @@ def test_partition_csv_rejects_a_wide_first_line_before_pandas_reads_it( @pytest.mark.parametrize("from_file", [False, True]) @pytest.mark.parametrize( - ("content", "n_cells"), + ("content", "n_measured_cells"), [ - # -- the context's restricted sniffer gives up on this one and Pandas sniffs it itself -- + # -- the context's restricted sniffer gives up on this one, so the delimiter is sniffed + # -- from the first line without restricting the candidates -- ("h" + "," * 99 + "\n" + "a\n" * 99, 100 * 100), - ("a;b;c\n1;2\n\n4\n", 3 * 3), # -- blank lines are not counted -- - # -- single-column file; Pandas sniffs "a" as the delimiter and reads 4 x 2 cells -- + # -- every line ending counts as a row, so the blank line makes this 4 rows (Pandas + # -- reads 3): the measure is an upper bound -- + ("a;b;c\n1;2\n\n4\n", 3 * 4), + # -- single-column file; the sniffer finds "a" as the delimiter, giving 2 columns -- ("a\nb\nc\nd\n", 4 * 2), - ('"x,\ny",z\n1\n', 2 * 2), # -- quoted delimiter and newline are not counted -- + ('"x,\ny",z\n1\n', 2 * 2), # -- quoted delimiter and newline start no field or row -- + ("a,b\r\nc,d\re,f", 3 * 2), # -- "\r\n", "\r" and an unterminated last line -- ], ) def test_partition_csv_limits_the_cells_the_file_spans( - content: str, n_cells: int, from_file: bool, tmp_path: Path, monkeypatch: pytest.MonkeyPatch + content: str, + n_measured_cells: int, + from_file: bool, + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, ): file_path = tmp_path / "table.csv" - file_path.write_text(content) + file_path.write_bytes(content.encode()) def partition(): if from_file: @@ -256,14 +264,52 @@ def partition(): return partition_csv(file=f) return partition_csv(str(file_path)) - monkeypatch.setenv("CSV_MAX_CELLS", str(n_cells)) + monkeypatch.setenv("CSV_MAX_CELLS", str(n_measured_cells)) assert len(partition()) == 1 - monkeypatch.setenv("CSV_MAX_CELLS", str(n_cells - 1)) + monkeypatch.setenv("CSV_MAX_CELLS", str(n_measured_cells - 1)) with pytest.raises(UnprocessableEntityError, match="CSV_MAX_CELLS"): partition() +@pytest.mark.parametrize( + "first_line", + [ + " \t\n", # -- a whitespace-only line is skipped, so the wide line is the first record -- + '"x\n"', # -- a quoted newline does not end the first record -- + "\x0c", # -- form-feed is not a line ending, so the commas are on the first line -- + 'a"b', # -- a quote inside a field does not open a quoted field -- + ], +) +def test_check_cell_count_measures_the_record_pandas_sizes_the_data_frame_by( + first_line: str, monkeypatch: pytest.MonkeyPatch +): + # -- Pandas reads 4 rows x 101 columns from each of these -- + file = io.BytesIO((first_line + "," * 100 + "\n" + "a,b\n" * 3).encode()) + monkeypatch.setenv("CSV_MAX_CELLS", "403") + + with pytest.raises(UnprocessableEntityError, match="101 columns"): + check_cell_count(file, ",", None) + + +def test_partition_csv_is_not_tricked_into_millions_of_rows_by_a_carriage_return(): + # -- Pandas 2.x's C tokenizer reads this as 262,145 rows -- + elements = partition_csv(file=io.BytesIO(b"a,b\n\r ,c\n")) + + assert elements[0].metadata.text_as_html == ( + "
ab
c
" + ) + + +def test_partition_csv_reads_a_file_with_no_usable_delimiter_as_one_column(): + # -- the sniffer picks the quote as the delimiter, which cannot delimit fields -- + elements = partition_csv(file=io.BytesIO(b'"a"\n"b"\n')) + + assert ( + elements[0].metadata.text_as_html == "
a
b
" + ) + + # ================================================================================================ # UNIT-TESTS # ================================================================================================ diff --git a/test_unstructured/partition/test_tsv.py b/test_unstructured/partition/test_tsv.py index 8eabe8027b..a46ee675d3 100644 --- a/test_unstructured/partition/test_tsv.py +++ b/test_unstructured/partition/test_tsv.py @@ -178,7 +178,7 @@ def test_partition_tsv_rejects_a_wide_first_line_before_pandas_reads_it( file_path.write_text("h" + "\t" * 4999 + "\n" + "a\n" * 5000) read_csv_ = mocker.patch.object(pd, "read_csv") - with pytest.raises(UnprocessableEntityError, match="1,001 rows x 5,000 columns"): + with pytest.raises(UnprocessableEntityError, match="rows x 5,000 columns"): if from_file: with open(file_path, "rb") as f: partition_tsv(file=f) @@ -203,3 +203,13 @@ def test_partition_tsv_partitions_a_file_at_the_cell_limit( elements = partition_tsv(str(file_path)) assert [e.text for e in elements] == ["a b c 1 2"] + + +def test_partition_tsv_reads_a_field_larger_than_the_csv_module_field_limit(tmp_path: Path): + # -- the `csv` module's default field limit is 128 KiB; Pandas has none -- + file_path = tmp_path / "big-field.tsv" + file_path.write_text("a\t" + "x" * 200_000 + "\n") + + (table,) = partition_tsv(str(file_path)) + + assert table.text == "a " + "x" * 200_000 diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 8d566c51b4..5607601634 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -3,9 +3,11 @@ import codecs import contextlib import csv +import io import itertools +import re from functools import cached_property -from typing import IO, Any, Iterator +from typing import IO, Any, Callable, Iterator, NoReturn import pandas as pd @@ -66,11 +68,9 @@ def partition_csv( with ctx.open() as file: check_cell_count(file, ctx.delimiter, ctx.encoding) with ctx.open() as file: - read_kw: dict = {"header": ctx.header, "sep": ctx.delimiter, "encoding": ctx.encoding} - # sep=None is not supported by the C engine; use Python engine to avoid ParserWarning. - if ctx.delimiter is None: - read_kw["engine"] = "python" - dataframe = pd.read_csv(file, **read_kw) + dataframe = read_delimited_text( + file, sep=ctx.delimiter, header=ctx.header, encoding=ctx.encoding + ) html_table = HtmlTable.from_html_text( dataframe.to_html(index=False, header=include_header, na_rep="") @@ -86,43 +86,205 @@ def partition_csv( return [Table(text=html_table.text, metadata=metadata, detection_origin=DETECTION_ORIGIN)] +def read_delimited_text( + file: IO[bytes], *, sep: str | None, header: int | None, encoding: str | None +) -> pd.DataFrame: + """Read delimited text from `file` into a data-frame, with its line endings normalized. + + Pandas 2.x's C tokenizer mishandles a "\r" line ending followed by a whitespace-only line + while skipping blank lines: it emits 2^18 empty rows for each, so a few bytes become millions + of rows inside `pd.read_csv()`, before any limit can be checked. Every line ending is therefore + converted to "\n" (including inside quoted fields) and given to Pandas as its only terminator. + + When `sep` is `None` the delimiter is sniffed from the first line by `_sniff_delimiter()`, the + same way `check_cell_count()` measures it, rather than leaving Pandas to sniff its own. + """ + encoding = encoding or "utf-8" + # -- like Pandas, drop a UTF-8 byte-order mark -- + if codecs.lookup(encoding).name == "utf-8": + encoding = "utf-8-sig" + text = _normalize_line_endings(file.read().decode(encoding)) + + if sep is None: + sep = _sniff_delimiter(text[: text.find("\n") + 1] if "\n" in text else text) + if sep is None: + # -- no usable delimiter, so the file is one column; split on a character it lacks -- + sep = next((c for c in _ABSENT_DELIMITER_CANDIDATES if c not in text), None) + if sep is None: + raise UnprocessableEntityError("Could not determine the delimiter of the file.") + # -- the C engine does not sniff and the Python engine needs no custom line terminator, + # -- so keep the Python engine this path has always used -- + return pd.read_csv(io.StringIO(text), sep=sep, header=header, engine="python") + return pd.read_csv(io.StringIO(text), sep=sep, header=header, lineterminator="\n") + + def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | None) -> None: """Raise `UnprocessableEntityError` when `file` would span more than `CSV_MAX_CELLS` cells. Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a tiny file whose first line is a long run of delimiters can span millions of cells. - The span is measured here by streaming the records, before Pandas allocates anything, and the - scan stops as soon as the limit is passed. `file` is read from its current position. + The span is measured here before Pandas allocates anything, reading `file` in fixed-size chunks + from its current position, with line endings normalized as `read_delimited_text()` does, and + the scan stops as soon as the limit is passed: + + - columns are the fields of the first record, counted with Pandas' quoting rules without + building the fields; + - rows are counted as line terminators after it. That is an upper bound, since blank lines and + newlines inside quoted fields count too, so the scan can over-count but never under-count. """ max_cells = env_config.CSV_MAX_CELLS - # -- `errors="replace"` so the scan never fails on a file Pandas would read -- - lines: Iterator[str] = codecs.getreader(encoding or "utf-8")(file, errors="replace") + chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding or "utf-8")) + def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: + raise UnprocessableEntityError( + f"File exceeds the maximum of {max_cells:,} table cells (CSV_MAX_CELLS): it spans" + f" at least {n_rows:,} rows x {n_cols:,} columns." + ) + + # -- with no delimiter given, `read_delimited_text()` sniffs one and uses Pandas' Python + # -- engine, which skips more kinds of blank line -- + python_engine = delimiter is None if delimiter is None: - # -- With `sep=None` Pandas sniffs the delimiter itself, from the first non-blank line - # -- and without restricting the candidates, so measure with the delimiter it will use. - first_line = next((line for line in lines if line.strip()), "") - try: - delimiter = csv.Sniffer().sniff(first_line).delimiter - except csv.Error: - delimiter = None - lines = itertools.chain([first_line], lines) - - # -- a single-column file has no delimiter; any character absent from each line will do -- - records = csv.reader(lines, delimiter=delimiter or "\n") - - n_rows, n_cols = 0, 0 - for record in records: - if not record: # -- Pandas skips blank lines -- - continue - if n_rows == 0: - n_cols = len(record) - n_rows += 1 - if n_rows * n_cols > max_cells: - raise UnprocessableEntityError( - f"File exceeds the maximum of {max_cells:,} table cells (CSV_MAX_CELLS): it spans" - f" at least {n_rows:,} rows x {n_cols:,} columns." - ) + first_line, chunks = _split_first_line(chunks) + # -- `None` when the file is read as one column -- + delimiter = _sniff_delimiter(first_line) + chunks = itertools.chain([first_line], chunks) + + n_cols, chunks = _first_record_width( + chunks, delimiter, python_engine, max_cells, raise_limit_exceeded + ) + if n_cols == 0: + return + + # -- the first record is row 1; each later "\n" ends a row, and so does end-of-file when the + # -- last line is unterminated -- + n_rows, unterminated = 1, False + for chunk in chunks: + if n_terminators := chunk.count("\n"): + n_rows += n_terminators + unterminated = not chunk.endswith("\n") + elif chunk: + unterminated = True + if (n_rows + unterminated) * n_cols > max_cells: + raise_limit_exceeded(n_rows + unterminated, n_cols) + + +# -- characters `read_delimited_text()` may split a one-column file on, if absent from it -- +_ABSENT_DELIMITER_CANDIDATES = "\x1f\x1e\x1d\x1c\x07\x08" + +_CSV_CHUNK_CHARS = 1 << 16 + + +def _normalize_line_endings(text: str) -> str: + """`text` with each "\r\n" and "\r" line ending replaced by "\n". + + Applied chunk by chunk, a "\r\n" split across chunks becomes two line endings, which only adds + a blank line. + """ + return text.replace("\r\n", "\n").replace("\r", "\n") + + +def _sniff_delimiter(first_line: str) -> str | None: + """The delimiter `csv.Sniffer` finds in `first_line`, or `None` when there is no usable one.""" + try: + delimiter = csv.Sniffer().sniff(first_line).delimiter + except csv.Error: + return None + # -- a quote or line ending cannot delimit fields -- + return None if delimiter in ('"', "\n", "\r") else delimiter + + +def _iter_decoded_chunks(file: IO[bytes], encoding: str) -> Iterator[str]: + """Generate the text of `file` in chunks of up to `_CSV_CHUNK_CHARS` bytes, decoded.""" + # -- `errors="replace"` so the scan never fails on a file Pandas would read -- + decoder = codecs.getincrementaldecoder(encoding)(errors="replace") + while chunk := file.read(_CSV_CHUNK_CHARS): + if text := decoder.decode(chunk): + yield text + if text := decoder.decode(b"", final=True): + yield text + + +def _split_first_line(chunks: Iterator[str]) -> tuple[str, Iterator[str]]: + """The first line of `chunks` with its "\n", and the chunks that follow it.""" + pending = "" + for chunk in chunks: + pending += chunk + if (end := pending.find("\n") + 1) > 0: + return pending[:end], itertools.chain([pending[end:]], chunks) + return pending, iter(()) + + +def _first_record_width( + chunks: Iterator[str], + delimiter: str | None, + python_engine: bool, + max_cells: int, + raise_limit_exceeded: Callable[[int, int], NoReturn], +) -> tuple[int, Iterator[str]]: + """The field count of the first non-blank record, and the chunks that follow it. + + Follows the quoting rules Pandas applies by default: a `"` opens a quoted field only at the + start of a field, `""` inside a quoted field is a literal quote, and delimiters and newlines + inside a quoted field are part of the field. A record with no delimiter or quote and only + whitespace is blank, as Pandas skips it: spaces and tabs for its C engine, anything + `str.isspace()` for its Python engine. A `None` delimiter means a one-column file. The count is + 0 when there is no non-blank record. The limit is applied to the record alone. + """ + is_blank = str.isspace if python_engine else (lambda text: not text.strip(" \t")) + special = re.compile(f"[{re.escape(delimiter)}\n]" if delimiter else "\n") + in_quotes = quote_pending = has_content = False + at_field_start = True + n_delimiters = 0 + + for chunk in chunks: + i, n = 0, len(chunk) + while i < n: + if in_quotes: + if quote_pending: + quote_pending = False + if chunk[i] == '"': # -- `""` is a literal quote -- + i += 1 + else: # -- the quote closed the quoted part; the field continues -- + in_quotes = False + continue + j = chunk.find('"', i) + if j < 0: + break + quote_pending, i = True, j + 1 + continue + + if at_field_start and chunk[i] == '"': + in_quotes = has_content = True + at_field_start = False + i += 1 + continue + + match = special.search(chunk, i) + j = match.start() if match else n + if j > i: + at_field_start = False + has_content = has_content or not is_blank(chunk[i:j]) + if match is None: + break + i = j + 1 + if chunk[j] == delimiter: + n_delimiters += 1 + has_content = at_field_start = True + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + elif has_content: # -- end of the first record -- + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + return n_delimiters + 1, itertools.chain([chunk[i:]], chunks) + else: # -- a blank line, which Pandas skips -- + at_field_start = True + + if not has_content: + return 0, iter(()) + if n_delimiters + 1 > max_cells: + raise_limit_exceeded(1, n_delimiters + 1) + return n_delimiters + 1, iter(()) class _CsvPartitioningContext: diff --git a/unstructured/partition/tsv.py b/unstructured/partition/tsv.py index 04643772be..a7463be4bf 100644 --- a/unstructured/partition/tsv.py +++ b/unstructured/partition/tsv.py @@ -2,8 +2,6 @@ from typing import IO, Any, Optional -import pandas as pd - from unstructured.chunking import add_chunking_strategy from unstructured.common.html_table import HtmlTable from unstructured.documents.elements import Element, ElementMetadata, Table @@ -13,7 +11,7 @@ spooled_to_bytes_io_if_needed, ) from unstructured.partition.common.metadata import apply_metadata, get_last_modified_date -from unstructured.partition.csv import check_cell_count +from unstructured.partition.csv import check_cell_count, read_delimited_text from unstructured.telemetry import partition_runtime_telemetry DETECTION_ORIGIN: str = "tsv" @@ -47,7 +45,8 @@ def partition_tsv( if filename: with open(filename, "rb") as f: check_cell_count(f, "\t", None) - dataframe = pd.read_csv(filename, sep="\t", header=header) + f.seek(0) + dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) else: assert file is not None # -- Note(scanny): `SpooledTemporaryFile` on Python<3.11 does not implement `.readable()` @@ -56,7 +55,7 @@ def partition_tsv( start = f.tell() check_cell_count(f, "\t", None) f.seek(start) - dataframe = pd.read_csv(f, sep="\t", header=header) + dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) html_table = HtmlTable.from_html_text( dataframe.to_html(index=False, header=include_header, na_rep="") From 6528f2f30ad7858ff49e40d7b8dc5bf180dee9c5 Mon Sep 17 00:00:00 2001 From: Yao You Date: Thu, 1 Oct 2026 21:04:34 -0500 Subject: [PATCH 3/6] fix(csv): convert only a lone "\r", keeping "\r\n" inside quoted fields Normalizing every line ending to "\n" also rewrote "\r\n" inside quoted fields, changing cell text from CRLF files (e.g. multi-line cells in exports). Only a lone "\r" triggers Pandas 2.x's row explosion, so convert just that and let Pandas read "\r\n" natively. The delimiter is still sniffed from the first line with its ending normalized, the same text the size check sniffs. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- unstructured/partition/csv.py | 26 +++++++++++++++----------- 2 files changed, 16 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index a736a74683..5a5dfd1613 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,7 @@ ### Fixes -- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. Line endings are also normalized to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read a `"\r"` line ending followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. When the delimiter is sniffed, it is now sniffed once and passed to Pandas, and a file with no usable delimiter is read as one column. +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. When the delimiter is sniffed, it is now sniffed once and passed to Pandas, and a file with no usable delimiter is read as one column. ## 0.27.11 diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 5607601634..526470173c 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -89,12 +89,13 @@ def partition_csv( def read_delimited_text( file: IO[bytes], *, sep: str | None, header: int | None, encoding: str | None ) -> pd.DataFrame: - """Read delimited text from `file` into a data-frame, with its line endings normalized. + """Read delimited text from `file` into a data-frame, with each lone "\r" made a "\n". - Pandas 2.x's C tokenizer mishandles a "\r" line ending followed by a whitespace-only line + Pandas 2.x's C tokenizer mishandles a lone "\r" line ending followed by a whitespace-only line while skipping blank lines: it emits 2^18 empty rows for each, so a few bytes become millions - of rows inside `pd.read_csv()`, before any limit can be checked. Every line ending is therefore - converted to "\n" (including inside quoted fields) and given to Pandas as its only terminator. + of rows inside `pd.read_csv()`, before any limit can be checked. A "\r" not followed by "\n" is + therefore converted to "\n" first. "\r\n" is left alone, so text inside quoted fields of a + "\r\n" file is unchanged. When `sep` is `None` the delimiter is sniffed from the first line by `_sniff_delimiter()`, the same way `check_cell_count()` measures it, rather than leaving Pandas to sniff its own. @@ -103,19 +104,22 @@ def read_delimited_text( # -- like Pandas, drop a UTF-8 byte-order mark -- if codecs.lookup(encoding).name == "utf-8": encoding = "utf-8-sig" - text = _normalize_line_endings(file.read().decode(encoding)) + text = _LONE_CARRIAGE_RETURN.sub("\n", file.read().decode(encoding)) if sep is None: - sep = _sniff_delimiter(text[: text.find("\n") + 1] if "\n" in text else text) + first_line = text[: text.find("\n") + 1] if "\n" in text else text + sep = _sniff_delimiter(_normalize_line_endings(first_line)) if sep is None: # -- no usable delimiter, so the file is one column; split on a character it lacks -- sep = next((c for c in _ABSENT_DELIMITER_CANDIDATES if c not in text), None) if sep is None: raise UnprocessableEntityError("Could not determine the delimiter of the file.") - # -- the C engine does not sniff and the Python engine needs no custom line terminator, - # -- so keep the Python engine this path has always used -- + # -- the C engine does not sniff, so keep the Python engine this path has always used -- return pd.read_csv(io.StringIO(text), sep=sep, header=header, engine="python") - return pd.read_csv(io.StringIO(text), sep=sep, header=header, lineterminator="\n") + return pd.read_csv(io.StringIO(text), sep=sep, header=header) + + +_LONE_CARRIAGE_RETURN = re.compile("\r(?!\n)") def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | None) -> None: @@ -124,8 +128,8 @@ def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | Non Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a tiny file whose first line is a long run of delimiters can span millions of cells. The span is measured here before Pandas allocates anything, reading `file` in fixed-size chunks - from its current position, with line endings normalized as `read_delimited_text()` does, and - the scan stops as soon as the limit is passed: + from its current position, with each line ending counted as one "\n", and the scan stops as + soon as the limit is passed: - columns are the fields of the first record, counted with Pandas' quoting rules without building the fields; From bdd02b8a81a660917d1974b6c9a64e60faa205a0 Mon Sep 17 00:00:00 2001 From: Yao You Date: Thu, 1 Oct 2026 21:32:10 -0500 Subject: [PATCH 4/6] fix(csv): stream the file to Pandas; match its BOM and blank-line rules Review follow-ups: - Stream the file to Pandas through a reader that decodes it in chunks and converts each lone "\r", holding back a "\r" that ends a chunk until the next one shows whether it starts a "\r\n". The file is no longer read into memory whole, decoded and copied before parsing. - Drop leading byte-order marks before both the size check and Pandas. The check decoded UTF-8 without dropping the BOM that the read dropped, so a BOM-prefixed blank line could hide a wide record from it. - Sniff the delimiter from the first non-blank line in both paths, so a leading blank line no longer turns a delimited file into one column. - For Pandas' Python engine, treat a record as blank when its one value, quoted or not, is whitespace, as that engine skips it. - Use the ASCII unit separator as the delimiter of a file with no usable one in both paths, so a streamed file needs no whole-text scan to pick one. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- test_unstructured/partition/test_csv.py | 51 +++++- unstructured/partition/csv.py | 201 ++++++++++++++++++------ 3 files changed, 200 insertions(+), 54 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5a5dfd1613..305422f12e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,7 @@ ### Fixes -- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. When the delimiter is sniffed, it is now sniffed once and passed to Pandas, and a file with no usable delimiter is read as one column. +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. The file is streamed to Pandas rather than read into memory whole. When the delimiter is sniffed, it is now sniffed once, from the first non-blank line, and passed to Pandas, and a file with no usable delimiter is read as one column. ## 0.27.11 diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index 140708931b..5e19d47340 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -31,7 +31,14 @@ from unstructured.cleaners.core import clean_extra_whitespace from unstructured.documents.elements import Table from unstructured.errors import UnprocessableEntityError -from unstructured.partition.csv import _CsvPartitioningContext, check_cell_count, partition_csv +from unstructured.partition import csv as csv_module +from unstructured.partition.csv import ( + _CsvPartitioningContext, + _first_record_width, + check_cell_count, + partition_csv, + read_delimited_text, +) from unstructured.partition.utils.constants import UNSTRUCTURED_INCLUDE_DEBUG_METADATA EXPECTED_FILETYPE = "text/csv" @@ -279,6 +286,7 @@ def partition(): '"x\n"', # -- a quoted newline does not end the first record -- "\x0c", # -- form-feed is not a line ending, so the commas are on the first line -- 'a"b', # -- a quote inside a field does not open a quoted field -- + "\ufeff\n", # -- a byte-order mark is dropped, leaving a blank line that is skipped -- ], ) def test_check_cell_count_measures_the_record_pandas_sizes_the_data_frame_by( @@ -301,6 +309,47 @@ def test_partition_csv_is_not_tricked_into_millions_of_rows_by_a_carriage_return ) +@pytest.mark.parametrize(("python_engine", "expected_width"), [(True, 11), (False, 1)]) +def test_first_record_width_applies_each_pandas_engines_blank_line_rule( + python_engine: bool, expected_width: int +): + # -- the Python engine skips a record whose one value is whitespace, even quoted; the C engine + # -- skips only an unquoted line of spaces and tabs -- + chunks = iter(['" "\n' + ";" * 10 + "\n"]) + + width, _ = _first_record_width(chunks, ";", python_engine, 10**9, Mock()) + + assert width == expected_width + + +def test_partition_csv_sniffs_the_delimiter_from_the_first_non_blank_line(): + # -- the context's sniffer only tries ",;|", so Pandas' delimiter is sniffed here -- + elements = partition_csv(file=io.BytesIO(b"\n\na\tb\tc\n1\t2\t3\n")) + + assert elements[0].metadata.text_as_html == ( + "" + "
abc
123
" + ) + + +def test_read_delimited_text_streams_the_file_in_chunks(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1) + read_sizes: list[int] = [] + + class RecordingBytesIO(io.BytesIO): + def read(self, size: int | None = -1) -> bytes: + read_sizes.append(-1 if size is None else size) + return super().read(size) + + # -- with 1-byte chunks every "\r" is held back to see whether a "\n" follows it -- + dataframe = read_delimited_text( + RecordingBytesIO(b'a,"x\r\ny"\r\nc,"d\re"\r'), sep=",", header=None, encoding=None + ) + + assert dataframe.values.tolist() == [["a", "x\r\ny"], ["c", "d\ne"]] + assert set(read_sizes) == {1} + + def test_partition_csv_reads_a_file_with_no_usable_delimiter_as_one_column(): # -- the sniffer picks the quote as the delimiter, which cannot delimit fields -- elements = partition_csv(file=io.BytesIO(b'"a"\n"b"\n')) diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 526470173c..77950f6570 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -93,33 +93,21 @@ def read_delimited_text( Pandas 2.x's C tokenizer mishandles a lone "\r" line ending followed by a whitespace-only line while skipping blank lines: it emits 2^18 empty rows for each, so a few bytes become millions - of rows inside `pd.read_csv()`, before any limit can be checked. A "\r" not followed by "\n" is - therefore converted to "\n" first. "\r\n" is left alone, so text inside quoted fields of a - "\r\n" file is unchanged. - - When `sep` is `None` the delimiter is sniffed from the first line by `_sniff_delimiter()`, the - same way `check_cell_count()` measures it, rather than leaving Pandas to sniff its own. + of rows inside `pd.read_csv()`, before any limit can be checked. Pandas therefore reads the + file through `_LoneCarriageReturnReader`, which streams it with each "\r" not followed by "\n" + converted to "\n". "\r\n" is left alone, so text inside quoted fields of a "\r\n" file is + unchanged. + + When `sep` is `None` the delimiter is sniffed by `_sniff_delimiter()` from the first non-blank + line, the same way `check_cell_count()` measures it, rather than leaving Pandas to sniff its + own. """ - encoding = encoding or "utf-8" - # -- like Pandas, drop a UTF-8 byte-order mark -- - if codecs.lookup(encoding).name == "utf-8": - encoding = "utf-8-sig" - text = _LONE_CARRIAGE_RETURN.sub("\n", file.read().decode(encoding)) - + reader = _LoneCarriageReturnReader(file, encoding) if sep is None: - first_line = text[: text.find("\n") + 1] if "\n" in text else text - sep = _sniff_delimiter(_normalize_line_endings(first_line)) - if sep is None: - # -- no usable delimiter, so the file is one column; split on a character it lacks -- - sep = next((c for c in _ABSENT_DELIMITER_CANDIDATES if c not in text), None) - if sep is None: - raise UnprocessableEntityError("Could not determine the delimiter of the file.") + sep = _sniff_delimiter(_normalize_line_endings(reader.peek_first_non_blank_line())) # -- the C engine does not sniff, so keep the Python engine this path has always used -- - return pd.read_csv(io.StringIO(text), sep=sep, header=header, engine="python") - return pd.read_csv(io.StringIO(text), sep=sep, header=header) - - -_LONE_CARRIAGE_RETURN = re.compile("\r(?!\n)") + return pd.read_csv(reader, sep=sep, header=header, engine="python") + return pd.read_csv(reader, sep=sep, header=header) def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | None) -> None: @@ -137,7 +125,7 @@ def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | Non newlines inside quoted fields count too, so the scan can over-count but never under-count. """ max_cells = env_config.CSV_MAX_CELLS - chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding or "utf-8")) + chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding)) def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: raise UnprocessableEntityError( @@ -149,10 +137,8 @@ def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: # -- engine, which skips more kinds of blank line -- python_engine = delimiter is None if delimiter is None: - first_line, chunks = _split_first_line(chunks) - # -- `None` when the file is read as one column -- + first_line, chunks = _peek_first_non_blank_line(chunks) delimiter = _sniff_delimiter(first_line) - chunks = itertools.chain([first_line], chunks) n_cols, chunks = _first_record_width( chunks, delimiter, python_engine, max_cells, raise_limit_exceeded @@ -173,12 +159,106 @@ def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: raise_limit_exceeded(n_rows + unterminated, n_cols) -# -- characters `read_delimited_text()` may split a one-column file on, if absent from it -- -_ABSENT_DELIMITER_CANDIDATES = "\x1f\x1e\x1d\x1c\x07\x08" +# -- the delimiter of a file with no usable one, which is then read as one column. It is the ASCII +# -- "unit separator", so a file that does contain it is split on it, the same way by both +# -- `check_cell_count()` and Pandas -- +_ONE_COLUMN_DELIMITER = "\x1f" _CSV_CHUNK_CHARS = 1 << 16 +class _LoneCarriageReturnReader(io.TextIOBase): + """Text stream of `file`, decoded in chunks, with each "\r" not followed by "\n" made a "\n". + + A "\r" ending a chunk is held back until the next chunk shows whether it starts a "\r\n". The + file is never held in memory whole; only the lines peeked by `.peek_first_non_blank_line()` + are buffered, until they are read. + """ + + def __init__(self, file: IO[bytes], encoding: str | None): + self._file = file + self._decoder = codecs.getincrementaldecoder(_python_encoding(encoding))() + self._buffer = "" + self._held_cr = False + self._at_start = True + self._eof = False + + def peek_first_non_blank_line(self) -> str: + """The first line that is not blank, or "" if there is none, without consuming it.""" + start = 0 + while True: + end = self._buffer.find("\n", start) + 1 + if end == 0: + if not self._fill(): + return self._buffer[start:] if self._buffer[start:].strip() else "" + continue + if self._buffer[start:end].strip(): + return self._buffer[start:end] + start = end + + def read(self, size: int | None = -1) -> str: + if size is None or size < 0: + while self._fill(): + pass + text, self._buffer = self._buffer, "" + return text + while len(self._buffer) < size and self._fill(): + pass + text, self._buffer = self._buffer[:size], self._buffer[size:] + return text + + def readline(self, size: int | None = -1) -> str: + while (end := self._buffer.find("\n") + 1) == 0 and self._fill(): + pass + if end == 0: + end = len(self._buffer) + if size is not None and 0 <= size < end: + end = size + text, self._buffer = self._buffer[:end], self._buffer[end:] + return text + + def readable(self) -> bool: + return True + + def _fill(self) -> bool: + """Append the next decoded chunk to the buffer; False when the file is exhausted.""" + if self._eof: + return False + chunk = self._file.read(_CSV_CHUNK_CHARS) + text = self._decoder.decode(chunk, final=not chunk) + if self._at_start: + text, self._at_start = _strip_leading_bom(text) + if self._held_cr: + text = "\r" + text + self._held_cr = False + if not chunk: + self._eof = True + elif text.endswith("\r"): + text, self._held_cr = text[:-1], True + self._buffer += _LONE_CARRIAGE_RETURN.sub("\n", text) + return True + + +_LONE_CARRIAGE_RETURN = re.compile("\r(?!\n)") + + +def _python_encoding(encoding: str | None) -> str: + """The codec to decode with: `encoding`, or UTF-8, dropping a UTF-8 byte-order mark as Pandas + does.""" + encoding = encoding or "utf-8" + return "utf-8-sig" if codecs.lookup(encoding).name == "utf-8" else encoding + + +def _strip_leading_bom(text: str) -> tuple[str, bool]: + """`text` without leading byte-order marks, and whether the start of the file is still ahead. + + Pandas drops a byte-order mark from the start of the text it is given, so every leading one is + dropped before Pandas or the scan sees the text, leaving none for Pandas to drop differently. + """ + text = text.lstrip("\ufeff") + return text, not text + + def _normalize_line_endings(text: str) -> str: """`text` with each "\r\n" and "\r" line ending replaced by "\n". @@ -188,40 +268,50 @@ def _normalize_line_endings(text: str) -> str: return text.replace("\r\n", "\n").replace("\r", "\n") -def _sniff_delimiter(first_line: str) -> str | None: - """The delimiter `csv.Sniffer` finds in `first_line`, or `None` when there is no usable one.""" +def _sniff_delimiter(first_line: str) -> str: + """The delimiter `csv.Sniffer` finds in `first_line`, or `_ONE_COLUMN_DELIMITER`.""" try: delimiter = csv.Sniffer().sniff(first_line).delimiter except csv.Error: - return None + return _ONE_COLUMN_DELIMITER # -- a quote or line ending cannot delimit fields -- - return None if delimiter in ('"', "\n", "\r") else delimiter + return _ONE_COLUMN_DELIMITER if delimiter in ('"', "\n", "\r") else delimiter -def _iter_decoded_chunks(file: IO[bytes], encoding: str) -> Iterator[str]: +def _iter_decoded_chunks(file: IO[bytes], encoding: str | None) -> Iterator[str]: """Generate the text of `file` in chunks of up to `_CSV_CHUNK_CHARS` bytes, decoded.""" # -- `errors="replace"` so the scan never fails on a file Pandas would read -- - decoder = codecs.getincrementaldecoder(encoding)(errors="replace") + decoder = codecs.getincrementaldecoder(_python_encoding(encoding))(errors="replace") + at_start = True while chunk := file.read(_CSV_CHUNK_CHARS): - if text := decoder.decode(chunk): + text = decoder.decode(chunk) + if at_start: + text, at_start = _strip_leading_bom(text) + if text: yield text if text := decoder.decode(b"", final=True): - yield text + yield text.lstrip("\ufeff") if at_start else text -def _split_first_line(chunks: Iterator[str]) -> tuple[str, Iterator[str]]: - """The first line of `chunks` with its "\n", and the chunks that follow it.""" - pending = "" +def _peek_first_non_blank_line(chunks: Iterator[str]) -> tuple[str, Iterator[str]]: + """The first line of `chunks` that is not blank, and all of `chunks`, that line included. + + The line ends with its "\n"; it is "" when every line is blank. + """ + pending, start = "", 0 for chunk in chunks: pending += chunk - if (end := pending.find("\n") + 1) > 0: - return pending[:end], itertools.chain([pending[end:]], chunks) - return pending, iter(()) + while (end := pending.find("\n", start) + 1) > 0: + if pending[start:end].strip(): + return pending[start:end], itertools.chain([pending], chunks) + start = end + line = pending[start:] if pending[start:].strip() else "" + return line, iter([pending]) def _first_record_width( chunks: Iterator[str], - delimiter: str | None, + delimiter: str, python_engine: bool, max_cells: int, raise_limit_exceeded: Callable[[int, int], NoReturn], @@ -230,13 +320,13 @@ def _first_record_width( Follows the quoting rules Pandas applies by default: a `"` opens a quoted field only at the start of a field, `""` inside a quoted field is a literal quote, and delimiters and newlines - inside a quoted field are part of the field. A record with no delimiter or quote and only - whitespace is blank, as Pandas skips it: spaces and tabs for its C engine, anything - `str.isspace()` for its Python engine. A `None` delimiter means a one-column file. The count is - 0 when there is no non-blank record. The limit is applied to the record alone. + inside a quoted field are part of the field. A record with no delimiter is blank when Pandas + skips it: for its C engine when it has no quote and only spaces and tabs; for its Python engine + when its one field's value, quoted or not, is empty or only `str.isspace()` characters. The + count is 0 when there is no non-blank record. The limit is applied to the record alone. """ is_blank = str.isspace if python_engine else (lambda text: not text.strip(" \t")) - special = re.compile(f"[{re.escape(delimiter)}\n]" if delimiter else "\n") + special = re.compile(f"[{re.escape(delimiter)}\n]") in_quotes = quote_pending = has_content = False at_field_start = True n_delimiters = 0 @@ -248,19 +338,26 @@ def _first_record_width( if quote_pending: quote_pending = False if chunk[i] == '"': # -- `""` is a literal quote -- + has_content = True i += 1 else: # -- the quote closed the quoted part; the field continues -- in_quotes = False continue j = chunk.find('"', i) + # -- the Python engine judges blankness by the value, quoted text included -- + quoted_text = chunk[i:] if j < 0 else chunk[i:j] + if quoted_text and not quoted_text.isspace(): + has_content = True if j < 0: break quote_pending, i = True, j + 1 continue if at_field_start and chunk[i] == '"': - in_quotes = has_content = True - at_field_start = False + # -- a quote alone makes a record non-blank for the C engine, which skips blank + # -- lines before parsing quotes -- + in_quotes, at_field_start = True, False + has_content = has_content or not python_engine i += 1 continue From b526efa29ef9c81e785d08fa5c42d7fd90ff453b Mon Sep 17 00:00:00 2001 From: Yao You Date: Fri, 2 Oct 2026 11:07:13 -0500 Subject: [PATCH 5/6] fix(csv): read the stream in linear time; restore TSV input handling Review follow-ups: - The reader kept text in one string and sliced the rest off on every read, so replaying a long blank prefix kept by the delimiter lookahead copied it once per line: quadratic in the prefix. Queue decoded chunks with a read offset instead, so each character is copied a bounded number of times however the text is read. The size check's own lookahead re-yields the original chunks rather than one grown string. - Count one more column when the first record is a header: Pandas takes a data field beyond the header's width as an implicit index. - partition_tsv() again decompresses a compressed filename (e.g. .tsv.gz), through Pandas' handle, the same content for the size check and the read; and a stream that cannot seek is spooled to a temporary file once rather than failing. - Cover chunk-boundary independence (multibyte, BOM, quotes, CRLF and lone "\r"), stopping to read once the limit is passed, and the replayed-blank-prefix fallback. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- test_unstructured/partition/test_csv.py | 123 +++++++++++++++++++ test_unstructured/partition/test_tsv.py | 41 +++++++ unstructured/partition/csv.py | 156 ++++++++++++++++++------ unstructured/partition/tsv.py | 54 ++++++-- 5 files changed, 328 insertions(+), 48 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 305422f12e..cbcfd59017 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,7 @@ ### Fixes -- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. The file is streamed to Pandas rather than read into memory whole. When the delimiter is sniffed, it is now sniffed once, from the first non-blank line, and passed to Pandas, and a file with no usable delimiter is read as one column. +- **Reject CSV and TSV files that span too many cells instead of exhausting memory.** Pandas sizes the data-frame by the first record and pads every shorter record out to that width, so a few-KB file whose first line is a long run of delimiters could span millions of cells and use many GB in `partition_csv()` and `partition_tsv()`. The file's span is now measured by streaming it before Pandas reads it, and a file spanning more than `CSV_MAX_CELLS` cells (default 5,000,000) raises `UnprocessableEntityError`. A lone `"\r"` line ending is also converted to `"\n"` before Pandas reads the file: Pandas 2.x's C tokenizer read one followed by a whitespace-only line as 2^18 empty rows, so 5 bytes became 262,145 rows. `"\r\n"` is left as is. The file is streamed to Pandas rather than read into memory whole, and `partition_tsv()` still decompresses a compressed filename (e.g. `.tsv.gz`) and accepts a stream that cannot seek. When the delimiter is sniffed, it is now sniffed once, from the first non-blank line, and passed to Pandas, and a file with no usable delimiter is read as one column. ## 0.27.11 diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index 5e19d47340..6714acef73 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -475,3 +475,126 @@ def it_raises_when_neither_file_path_nor_file_is_provided(self): @pytest.fixture() def get_last_modified_date_(self, request: FixtureRequest) -> Mock: return function_mock(request, "unstructured.partition.csv.get_last_modified_date") + + +# -- streaming and boundaries -------------------------------------------------------------------- + + +def test_partition_csv_reads_past_a_long_blank_prefix_in_linear_time( + monkeypatch: pytest.MonkeyPatch, +): + # -- no delimiter in the context's sample, so the delimiter is sniffed past the blank lines and + # -- every blank line is replayed to Pandas' Python engine -- + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1024) + data = b"\n" * 200_000 + b"a\tb\n1\t2\n" + + elements = partition_csv(file=io.BytesIO(data)) + + assert elements[0].metadata.text_as_html == ( + "
ab
12
" + ) + + +def test_lone_carriage_return_reader_copies_each_character_a_bounded_number_of_times( + monkeypatch: pytest.MonkeyPatch, +): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1024) + data = b"\n" * 200_000 + b"a\tb\n1\t2\n" + reader = csv_module._LoneCarriageReturnReader(io.BytesIO(data), None) + + assert reader.peek_first_non_blank_line() == "a\tb\n" + lines = list(iter(reader.readline, "")) + + assert len(lines) == 200_002 + assert "".join(lines) == data.decode() + assert reader._chars_copied <= 3 * len(data) + + +def test_peek_first_non_blank_line_replays_the_original_chunks(): + chunks = ["\n" * 10, " \t\n", "a,", "b\n", "c,d\n"] + + line, replay = csv_module._peek_first_non_blank_line(iter(chunks)) + + assert line == "a,b\n" + assert all(a is b for a, b in zip(replay, chunks, strict=True)) + + +def test_partition_csv_counts_the_implicit_index_column_of_a_header( + monkeypatch: pytest.MonkeyPatch, +): + # -- data rows one field wider than the header: Pandas uses the extra field as the index -- + data = b"a,b\n1,2,3\n4,5,6\n" + monkeypatch.setenv("CSV_MAX_CELLS", str(3 * 3)) + assert len(partition_csv(file=io.BytesIO(data), include_header=True)) == 1 + + monkeypatch.setenv("CSV_MAX_CELLS", str(3 * 3 - 1)) + with pytest.raises(UnprocessableEntityError, match="3 columns"): + partition_csv(file=io.BytesIO(data), include_header=True) + + +@pytest.mark.parametrize( + "data", + [ + 'é,ü\r\nx,"a\r\nb"\r\n', # -- BOM, multibyte, CRLF inside and outside quotes -- + '日本,語\r1,"2\r3"\r\r\n', # -- lone "\r" inside and outside quotes -- + '\n"q""uote","é"\n', # -- repeated BOM, doubled quote -- + ], +) +@pytest.mark.parametrize("chunk_size", [1, 3]) +def test_reading_does_not_depend_on_where_chunks_split_the_file( + data: str, chunk_size: int, monkeypatch: pytest.MonkeyPatch +): + encoded = data.encode() + expected_text = csv_module._LoneCarriageReturnReader(io.BytesIO(encoded), None).read() + expected_frame = read_delimited_text(io.BytesIO(encoded), sep=",", header=None, encoding=None) + expected_cells = expected_frame.shape[0] * expected_frame.shape[1] + + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", chunk_size) + + assert csv_module._LoneCarriageReturnReader(io.BytesIO(encoded), None).read() == expected_text + frame = read_delimited_text(io.BytesIO(encoded), sep=",", header=None, encoding=None) + assert frame.equals(expected_frame) + monkeypatch.setenv("CSV_MAX_CELLS", str(expected_cells - 1)) + with pytest.raises(UnprocessableEntityError): + check_cell_count(io.BytesIO(encoded), ",", None) + + +@pytest.mark.parametrize( + "data", + [ + b"h" + + b"," * 100 + + b"\n" + + b"a\n" * 100_000, # -- the first record alone passes the limit -- + b"a,b\n" * 100_000, # -- the rows pass the limit -- + ], +) +def test_check_cell_count_stops_reading_once_the_limit_is_passed( + data: bytes, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 64) + monkeypatch.setenv("CSV_MAX_CELLS", "50") + n_bytes_read = 0 + + class CountingBytesIO(io.BytesIO): + def read(self, size: int | None = -1) -> bytes: + nonlocal n_bytes_read + chunk = super().read(size) + n_bytes_read += len(chunk) + return chunk + + with pytest.raises(UnprocessableEntityError): + check_cell_count(CountingBytesIO(data), ",", None) + + assert n_bytes_read <= 4 * 64 + + +def test_peek_first_non_blank_line_completes_a_line_across_a_chunk_of_only_a_held_carriage_return( + monkeypatch: pytest.MonkeyPatch, +): + # -- with 1-byte chunks the "\r" chunk queues nothing until the "\n" after it arrives -- + monkeypatch.setattr(csv_module, "_CSV_CHUNK_CHARS", 1) + reader = csv_module._LoneCarriageReturnReader(io.BytesIO(b"a\r\nb\n"), None) + + assert reader.peek_first_non_blank_line() == "a\r\n" + assert reader.read() == "a\r\nb\n" diff --git a/test_unstructured/partition/test_tsv.py b/test_unstructured/partition/test_tsv.py index a46ee675d3..98f11a304b 100644 --- a/test_unstructured/partition/test_tsv.py +++ b/test_unstructured/partition/test_tsv.py @@ -2,7 +2,10 @@ from __future__ import annotations +import gzip +import io from pathlib import Path +from typing import Any import pandas as pd import pytest @@ -213,3 +216,41 @@ def test_partition_tsv_reads_a_field_larger_than_the_csv_module_field_limit(tmp_ (table,) = partition_tsv(str(file_path)) assert table.text == "a " + "x" * 200_000 + + +def test_partition_tsv_decompresses_a_compressed_filename( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + file_path = tmp_path / "table.tsv.gz" + with gzip.open(file_path, "wb") as f: + f.write(b"a\tb\n1\t2\n") + + (table,) = partition_tsv(str(file_path)) + + assert table.text == "a b 1 2" + + # -- the size check measures the decompressed content -- + monkeypatch.setenv("CSV_MAX_CELLS", "3") + with pytest.raises(UnprocessableEntityError, match="CSV_MAX_CELLS"): + partition_tsv(str(file_path)) + + +def test_partition_tsv_reads_a_stream_that_cannot_seek(): + class Pipe(io.RawIOBase): + def __init__(self, data: bytes): + self._data = io.BytesIO(data) + + def readable(self) -> bool: + return True + + def seekable(self) -> bool: + return False + + def readinto(self, buffer: Any) -> int: + chunk = self._data.read(len(buffer)) + buffer[: len(chunk)] = chunk + return len(chunk) + + (table,) = partition_tsv(file=io.BufferedReader(Pipe(b"a\tb\n1\t2\n"))) + + assert table.text == "a b 1 2" diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 77950f6570..1b34f6a0a9 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -1,6 +1,7 @@ from __future__ import annotations import codecs +import collections import contextlib import csv import io @@ -66,7 +67,7 @@ def partition_csv( csv.field_size_limit(CSV_FIELD_LIMIT) with ctx.open() as file: - check_cell_count(file, ctx.delimiter, ctx.encoding) + check_cell_count(file, ctx.delimiter, ctx.encoding, header=ctx.header is not None) with ctx.open() as file: dataframe = read_delimited_text( file, sep=ctx.delimiter, header=ctx.header, encoding=ctx.encoding @@ -110,7 +111,9 @@ def read_delimited_text( return pd.read_csv(reader, sep=sep, header=header) -def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | None) -> None: +def check_cell_count( + file: IO[bytes], delimiter: str | None, encoding: str | None, *, header: bool = False +) -> None: """Raise `UnprocessableEntityError` when `file` would span more than `CSV_MAX_CELLS` cells. Pandas sizes the data-frame by the first record and pads every shorter record out to that @@ -123,6 +126,9 @@ def check_cell_count(file: IO[bytes], delimiter: str | None, encoding: str | Non building the fields; - rows are counted as line terminators after it. That is an upper bound, since blank lines and newlines inside quoted fields count too, so the scan can over-count but never under-count. + + With `header`, the first record is the header, and Pandas accepts data rows one field wider + than it (using the extra field as an implicit index), so one more column is counted. """ max_cells = env_config.CSV_MAX_CELLS chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding)) @@ -145,6 +151,7 @@ def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: ) if n_cols == 0: return + n_cols += header # -- the first record is row 1; each later "\n" ends a row, and so does end-of-file when the # -- last line is unterminated -- @@ -171,57 +178,99 @@ class _LoneCarriageReturnReader(io.TextIOBase): """Text stream of `file`, decoded in chunks, with each "\r" not followed by "\n" made a "\n". A "\r" ending a chunk is held back until the next chunk shows whether it starts a "\r\n". The - file is never held in memory whole; only the lines peeked by `.peek_first_non_blank_line()` - are buffered, until they are read. + file is never held in memory whole: decoded chunks wait in a queue until they are read, and + only the lines `.peek_first_non_blank_line()` looks past stay queued after it returns. + + Each character is copied a bounded number of times however the text is read, so reading is + linear in the size of the file; `._chars_copied` counts the copies. """ def __init__(self, file: IO[bytes], encoding: str | None): self._file = file self._decoder = codecs.getincrementaldecoder(_python_encoding(encoding))() - self._buffer = "" + self._chunks: collections.deque[str] = collections.deque() + self._offset = 0 # -- read position in the first queued chunk -- self._held_cr = False self._at_start = True self._eof = False + self._chars_copied = 0 def peek_first_non_blank_line(self) -> str: """The first line that is not blank, or "" if there is none, without consuming it.""" - start = 0 + line_parts: list[str] = [] # -- the current line, kept only until it proves blank -- + idx, start = 0, self._offset while True: - end = self._buffer.find("\n", start) + 1 - if end == 0: + if idx == len(self._chunks): if not self._fill(): - return self._buffer[start:] if self._buffer[start:].strip() else "" + line = self._join(line_parts) + return line if line.strip() else "" continue - if self._buffer[start:end].strip(): - return self._buffer[start:end] - start = end + chunk = self._chunks[idx] + end = chunk.find("\n", start) + 1 or len(chunk) + line_parts.append(self._slice(chunk, start, end)) + if not line_parts[-1].isspace() and line_parts[-1]: + # -- not blank; complete the line, then return it -- + while not line_parts[-1].endswith("\n"): + idx, start = idx + 1, 0 + # -- a fill can queue nothing (e.g. a chunk holding only a held-back "\r") -- + while idx == len(self._chunks): + if not self._fill(): + return self._join(line_parts) + chunk = self._chunks[idx] + end = chunk.find("\n") + 1 or len(chunk) + line_parts.append(self._slice(chunk, 0, end)) + return self._join(line_parts) + if line_parts[-1].endswith("\n"): + line_parts = [] + if end == len(chunk): + idx, start = idx + 1, 0 + else: + start = end def read(self, size: int | None = -1) -> str: - if size is None or size < 0: - while self._fill(): - pass - text, self._buffer = self._buffer, "" - return text - while len(self._buffer) < size and self._fill(): - pass - text, self._buffer = self._buffer[:size], self._buffer[size:] - return text + parts: list[str] = [] + remaining = -1 if size is None or size < 0 else size + while remaining != 0 and (self._chunks or self._fill()): + if not self._chunks: + continue + chunk = self._chunks[0] + end = len(chunk) if remaining < 0 else min(len(chunk), self._offset + remaining) + parts.append(self._slice(chunk, self._offset, end)) + if remaining > 0: + remaining -= end - self._offset + self._advance(end) + return self._join(parts) def readline(self, size: int | None = -1) -> str: - while (end := self._buffer.find("\n") + 1) == 0 and self._fill(): - pass - if end == 0: - end = len(self._buffer) - if size is not None and 0 <= size < end: - end = size - text, self._buffer = self._buffer[:end], self._buffer[end:] - return text + parts: list[str] = [] + remaining = -1 if size is None or size < 0 else size + while remaining != 0 and (self._chunks or self._fill()): + if not self._chunks: + continue + chunk = self._chunks[0] + end = chunk.find("\n", self._offset) + 1 or len(chunk) + if remaining > 0: + end = min(end, self._offset + remaining) + remaining -= end - self._offset + parts.append(self._slice(chunk, self._offset, end)) + self._advance(end) + if parts[-1].endswith("\n"): + break + return self._join(parts) def readable(self) -> bool: return True + def _advance(self, end: int) -> None: + """Move the read position to `end` in the first queued chunk, dropping it once read.""" + if end == len(self._chunks[0]): + self._chunks.popleft() + self._offset = 0 + else: + self._offset = end + def _fill(self) -> bool: - """Append the next decoded chunk to the buffer; False when the file is exhausted.""" + """Queue the next decoded chunk; False when the file is exhausted.""" if self._eof: return False chunk = self._file.read(_CSV_CHUNK_CHARS) @@ -235,9 +284,22 @@ def _fill(self) -> bool: self._eof = True elif text.endswith("\r"): text, self._held_cr = text[:-1], True - self._buffer += _LONE_CARRIAGE_RETURN.sub("\n", text) + if text: + self._chunks.append(_LONE_CARRIAGE_RETURN.sub("\n", text)) return True + def _join(self, parts: list[str]) -> str: + if len(parts) == 1: + return parts[0] + self._chars_copied += sum(map(len, parts)) + return "".join(parts) + + def _slice(self, chunk: str, start: int, end: int) -> str: + if start == 0 and end == len(chunk): + return chunk + self._chars_copied += end - start + return chunk[start:end] + _LONE_CARRIAGE_RETURN = re.compile("\r(?!\n)") @@ -296,17 +358,33 @@ def _iter_decoded_chunks(file: IO[bytes], encoding: str | None) -> Iterator[str] def _peek_first_non_blank_line(chunks: Iterator[str]) -> tuple[str, Iterator[str]]: """The first line of `chunks` that is not blank, and all of `chunks`, that line included. - The line ends with its "\n"; it is "" when every line is blank. + The line ends with its "\n"; it is "" when every line is blank. Each chunk read is queued + unchanged for the returned chunks, and only the non-blank line itself is copied. """ - pending, start = "", 0 + read: list[str] = [] + line_parts: list[str] = [] # -- the current line, kept only until it proves blank -- for chunk in chunks: - pending += chunk - while (end := pending.find("\n", start) + 1) > 0: - if pending[start:end].strip(): - return pending[start:end], itertools.chain([pending], chunks) + read.append(chunk) + start = 0 + while start < len(chunk): + end = chunk.find("\n", start) + 1 or len(chunk) + segment = chunk[start:end] + line_parts.append(segment) + if segment and not segment.isspace(): + if not segment.endswith("\n"): + # -- complete the line from the chunks that follow -- + for more in chunks: + read.append(more) + end = more.find("\n") + 1 or len(more) + line_parts.append(more[:end]) + if line_parts[-1].endswith("\n"): + break + return "".join(line_parts), itertools.chain(read, chunks) + if segment.endswith("\n"): + line_parts = [] start = end - line = pending[start:] if pending[start:].strip() else "" - return line, iter([pending]) + line = "".join(line_parts) + return (line if line.strip() else ""), iter(read) def _first_record_width( diff --git a/unstructured/partition/tsv.py b/unstructured/partition/tsv.py index a7463be4bf..dd5c22911a 100644 --- a/unstructured/partition/tsv.py +++ b/unstructured/partition/tsv.py @@ -1,6 +1,11 @@ from __future__ import annotations -from typing import IO, Any, Optional +import contextlib +import shutil +import tempfile +from typing import IO, Any, Iterator, Optional, cast + +from pandas.io.common import get_handle from unstructured.chunking import add_chunking_strategy from unstructured.common.html_table import HtmlTable @@ -43,19 +48,22 @@ def partition_tsv( header = 0 if include_header else None if filename: - with open(filename, "rb") as f: - check_cell_count(f, "\t", None) - f.seek(0) + # -- like `pd.read_csv(filename)`, decompress a file named e.g. "x.tsv.gz"; the size check + # -- and the read each get their own decompressing handle on the same content -- + with _open_decompressed(filename) as f: + check_cell_count(f, "\t", None, header=include_header) + with _open_decompressed(filename) as f: dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) else: assert file is not None # -- Note(scanny): `SpooledTemporaryFile` on Python<3.11 does not implement `.readable()` # -- which triggers an exception on `pd.DataFrame.read_csv()` call. f = spooled_to_bytes_io_if_needed(file) - start = f.tell() - check_cell_count(f, "\t", None) - f.seek(start) - dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) + with _rereadable(f) as f: + start = f.tell() + check_cell_count(f, "\t", None, header=include_header) + f.seek(start) + dataframe = read_delimited_text(f, sep="\t", header=header, encoding=None) html_table = HtmlTable.from_html_text( dataframe.to_html(index=False, header=include_header, na_rep="") @@ -69,3 +77,33 @@ def partition_tsv( metadata.detection_origin = DETECTION_ORIGIN return [Table(text=html_table.text, metadata=metadata)] + + +@contextlib.contextmanager +def _open_decompressed(filename: str) -> Iterator[IO[bytes]]: + """Open `filename` for reading bytes, decompressed as its extension (e.g. ".gz") implies.""" + with get_handle(filename, "rb", compression="infer", is_text=False) as handles: + yield cast(IO[bytes], handles.handle) + + +@contextlib.contextmanager +def _rereadable(file: IO[bytes]) -> Iterator[IO[bytes]]: + """`file` itself if it can seek, otherwise a spooled copy of the rest of it. + + The size check and the read each read the file, so a stream that cannot seek back (e.g. a + pipe or socket) is copied once to a temporary file, which spills to disk when large. + """ + try: + seekable = file.seekable() + except (AttributeError, ValueError): + seekable = False + if seekable: + yield file + return + with tempfile.SpooledTemporaryFile(max_size=_SPOOL_MAX_MEMORY_BYTES) as copy: + shutil.copyfileobj(file, copy) + copy.seek(0) + yield cast(IO[bytes], copy) + + +_SPOOL_MAX_MEMORY_BYTES = 64 * 1024 * 1024 From b6ad778982cb4e0f1bd243e48d98eba424f269c5 Mon Sep 17 00:00:00 2001 From: Yao You Date: Fri, 2 Oct 2026 12:42:35 -0500 Subject: [PATCH 6/6] fix(csv): count every implicit-index column a header admits With a header, Pandas turns the fields a data row has beyond the header into implicit-index columns: its C engine takes their number from the first data row, and its Python engine also from the second, which can make the first data row the index names. Counting one extra column under-counted a short header followed by a very wide row. Measure the header and the first two data records with the same quote-aware scanner and count the widest; any later, wider row is an error in Pandas. Also describe the check as a conservative cell estimate rather than a memory or byte bound. Co-Authored-By: Claude Opus 5.5 (1M context) --- test_unstructured/partition/test_csv.py | 17 +++++++++++++ test_unstructured/partition/test_tsv.py | 32 +++++++++++++++++++++++++ unstructured/partition/csv.py | 30 ++++++++++++++++++----- 3 files changed, 73 insertions(+), 6 deletions(-) diff --git a/test_unstructured/partition/test_csv.py b/test_unstructured/partition/test_csv.py index 6714acef73..15c4c6b6d4 100644 --- a/test_unstructured/partition/test_csv.py +++ b/test_unstructured/partition/test_csv.py @@ -598,3 +598,20 @@ def test_peek_first_non_blank_line_completes_a_line_across_a_chunk_of_only_a_hel assert reader.peek_first_non_blank_line() == "a\r\n" assert reader.read() == "a\r\nb\n" + + +def test_partition_csv_counts_index_names_the_python_engine_infers( + mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + # -- the context's sniffer only tries ",;|", so this tab-delimited file is read by Pandas' + # -- Python engine. Its second data row is as wide as the first plus the header, so Pandas + # -- makes the first data row the index names: 3 columns, from a 2-field header -- + data = b"a\tb\nidx\ni\t1\t2\n" + b"j\t3\t4\n" * 10 + pandas_cells = 13 * 3 + monkeypatch.setenv("CSV_MAX_CELLS", str(pandas_cells - 1)) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="3 columns"): + partition_csv(file=io.BytesIO(data), include_header=True) + + read_csv_.assert_not_called() diff --git a/test_unstructured/partition/test_tsv.py b/test_unstructured/partition/test_tsv.py index 98f11a304b..59bf820f2a 100644 --- a/test_unstructured/partition/test_tsv.py +++ b/test_unstructured/partition/test_tsv.py @@ -20,6 +20,7 @@ ) from test_unstructured.unit_utils import assert_round_trips_through_JSON, example_doc_path from unstructured.chunking.title import chunk_by_title +from unstructured.common.html_table import HtmlTable from unstructured.documents.elements import Table from unstructured.errors import UnprocessableEntityError from unstructured.partition.tsv import partition_tsv @@ -254,3 +255,34 @@ def readinto(self, buffer: Any) -> int: (table,) = partition_tsv(file=io.BufferedReader(Pipe(b"a\tb\n1\t2\n"))) assert table.text == "a b 1 2" + + +def test_partition_tsv_counts_every_implicit_index_column_before_pandas_reads( + tmp_path: Path, mocker: MockFixture, monkeypatch: pytest.MonkeyPatch +): + # -- a 1-field header, a 100-field first data row (99 implicit-index columns) and a ragged + # -- tail: Pandas builds 1,002 rows x 100 columns -- + file_path = tmp_path / "index.tsv" + file_path.write_text("h\n" + "\t".join(f"v{i}" for i in range(100)) + "\n" + "a\n" * 1000) + monkeypatch.setenv("CSV_MAX_CELLS", str(1002 * 100 - 1)) + read_csv_ = mocker.patch.object(pd, "read_csv") + + with pytest.raises(UnprocessableEntityError, match="100 columns"): + partition_tsv(str(file_path), include_header=True) + + read_csv_.assert_not_called() + + +def test_partition_tsv_with_implicit_index_columns_matches_pandas_within_the_limit( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + monkeypatch.delenv("CSV_MAX_CELLS", raising=False) + file_path = tmp_path / "index.tsv" + file_path.write_text("h\n" + "x\ty\tz\n" + "a\n" * 3) + + (table,) = partition_tsv(str(file_path), include_header=True) + + expected = pd.read_csv(file_path, sep="\t", header=0).to_html( + index=False, header=True, na_rep="" + ) + assert table.text == HtmlTable.from_html_text(expected).text diff --git a/unstructured/partition/csv.py b/unstructured/partition/csv.py index 1b34f6a0a9..3c9f772167 100644 --- a/unstructured/partition/csv.py +++ b/unstructured/partition/csv.py @@ -127,8 +127,15 @@ def check_cell_count( - rows are counted as line terminators after it. That is an upper bound, since blank lines and newlines inside quoted fields count too, so the scan can over-count but never under-count. - With `header`, the first record is the header, and Pandas accepts data rows one field wider - than it (using the extra field as an implicit index), so one more column is counted. + With `header`, the first record is the header, and Pandas makes the fields a data row has + beyond it an implicit index: its C engine takes their number from the first data row, and its + Python engine also from the second, which can make the first data row the index names. So the + columns are the widest of the header and the first two data records, each measured the same + way; any later, wider row is an error in Pandas. + + This is a conservative estimate of the cells Pandas will build, not a bound on memory or bytes + read: blank lines held while sniffing, or a stream spooled because it cannot seek, still take + space in proportion to the input. """ max_cells = env_config.CSV_MAX_CELLS chunks = (_normalize_line_endings(c) for c in _iter_decoded_chunks(file, encoding)) @@ -151,11 +158,22 @@ def raise_limit_exceeded(n_rows: int, n_cols: int) -> NoReturn: ) if n_cols == 0: return - n_cols += header + n_rows = 1 - # -- the first record is row 1; each later "\n" ends a row, and so does end-of-file when the - # -- last line is unterminated -- - n_rows, unterminated = 1, False + if header: + # -- the first two data records decide how many implicit-index columns Pandas adds -- + for _ in range(2): + width, chunks = _first_record_width( + chunks, delimiter, python_engine, max_cells, raise_limit_exceeded + ) + if width == 0: + return + n_rows, n_cols = n_rows + 1, max(n_cols, width) + if n_rows * n_cols > max_cells: + raise_limit_exceeded(n_rows, n_cols) + + # -- each later "\n" ends a row, and so does end-of-file when the last line is unterminated -- + unterminated = False for chunk in chunks: if n_terminators := chunk.count("\n"): n_rows += n_terminators