From 24287e5777035df19650d5dbbb7dd88c603690a6 Mon Sep 17 00:00:00 2001 From: linhongyu510 Date: Thu, 24 Sep 2026 19:31:54 +0800 Subject: [PATCH 1/2] fix(xpacks/llm): stop TokenCountSplitter dropping characters at punctuation cuts TokenCountSplitter cut a chunk at the last punctuation mark, then advanced the token cursor by len(encode(kept_text)). Token boundaries need not align with the character cut, so the token straddling the cut was skipped and the characters between the punctuation and the next token boundary were lost. Example: TokenCountSplitter(min_tokens=1, max_tokens=3).chunk("a.b.c.d.e.f.g.h.i.j.k.l.") dropped "f", "i", "l" -- rejoining the chunks yielded "a.b.c.d.e..g.h..j.k.." . Advance instead by exactly the whole source tokens whose decoded text stays within the cut (the longest token prefix that is still a prefix of the kept text), and re-emit the straddling token in the next chunk. Concatenating the chunks now reconstructs the (unicode-normalized) input. Also guards the cursor against a zero-token advance that could stall the loop. --- CHANGELOG.md | 7 +++++++ python/pathway/xpacks/llm/splitters.py | 21 +++++++++++++++++-- .../xpacks/llm/tests/test_splitters.py | 14 +++++++++++++ 3 files changed, 40 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8142bc16f..35f384e9c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,13 @@ This project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.htm - The base installation no longer depends on `boto3`. Everything it backed works without it: the region of an S3 bucket is detected with a single anonymous request, and AWS credentials for the Delta Lake connector are resolved by the official AWS SDK chain built into the engine — environment variables, profile files (including SSO and assume-role), ECS and EC2 instance credentials. If your own code imports `boto3` and relied on Pathway pulling it in, add it to your project's dependencies. - Building a pipeline is now significantly faster: the location information attached to every expression and operator (used in the `Occurred here` part of error messages) is collected without materializing full stack traces. The speedup is most visible in code that constructs many small pipelines, such as unit test suites, where graph construction could previously dominate the run time — measured up to ~6x faster on pytest suites of small-data pipeline tests. Error messages are unchanged. +### Fixed +- `pw.io.rabbitmq.read` in the static mode no longer waits forever on a stream that holds no messages: such a read finishes with an empty table. The end of the stream is also waited for longer (15s instead of 5s) before the stream is believed empty, so a slow broker is less likely to turn a static read into an empty one. +- Output connectors with a `time` column (`pw.io.deltalake`, `pw.io.iceberg`, `pw.io.postgres`, `pw.io.mysql`, `pw.io.mssql`, `pw.io.sqlite`, `pw.io.duckdb`, `pw.io.mongodb`) no longer fail with a `worker panic: … TryFromIntError` (or an `IntOutOfRange` error) when they receive the rows that a `windowby`/`forget` buffer with a `delay` releases once the input ends. Such rows are written with `time` equal to the maximal 64-bit signed integer, marking the end of the stream. +- `pw.io.mqtt.read` with persistence enabled no longer loses messages on restart. With `qos` 1 or 2, a message is acknowledged to the broker only after a durable checkpoint covers it, and the broker session now survives restarts — so both the messages received shortly before a crash and the messages published while the pipeline was down are redelivered (at-least-once delivery). +- `pw.io.nats.read` with persistence enabled and a JetStream stream no longer loses the messages read between the last checkpoint and a crash: they are acknowledged only once a checkpoint covers them, so the server redelivers them after the restart (at-least-once delivery). +- `pathway.xpacks.llm.splitters.TokenCountSplitter` no longer drops characters when a chunk is cut at a punctuation mark. The cursor was advanced by the number of tokens the *kept* text re-encodes to, but token boundaries need not align with the character cut, so the token straddling the cut was skipped and the characters between the punctuation and the next token boundary were lost (e.g. `"a.b.c.d.e.f."` split with `max_tokens=3` dropped `f`). The splitter now advances by exactly the whole tokens kept in the chunk and re-emits the straddling token in the next chunk, so concatenating the chunks reconstructs the input. + ### Added - `pw.io.python.ConnectorSubject` now accepts `session_type="upsert"`: a row sent with an already-present primary key replaces the previous version instead of being an error, and the new public `delete` method removes a row. This is the intended mode for polling-style connectors that periodically re-send everything they see; it requires the schema to define a primary key. - `pw.io.postgres.write` and `pw.io.postgres.write_snapshot` now support PostGIS `geometry` and `geography` columns: a `str` column holding WKT (in any dimension, optionally with an `SRID=n;` prefix) or hex-encoded EWKB is converted to the binary EWKB form the server expects. `pw.io.postgres.read` reads such columns back as hex-encoded EWKB strings, so geometry values round-trip. diff --git a/python/pathway/xpacks/llm/splitters.py b/python/pathway/xpacks/llm/splitters.py index ea8457af7..3c754f1a5 100644 --- a/python/pathway/xpacks/llm/splitters.py +++ b/python/pathway/xpacks/llm/splitters.py @@ -259,6 +259,7 @@ def chunk(self, text: str, metadata: dict = {}, **kwargs) -> list[tuple[str, dic while i < len(tokens): chunk_tokens = tokens[i : i + max_tokens] chunk = tokenizer.decode(chunk_tokens) + consumed = len(chunk_tokens) last_punctuation = max( [chunk.rfind(p) for p in self.PUNCTUATION], default=-1 ) @@ -266,8 +267,24 @@ def chunk(self, text: str, metadata: dict = {}, **kwargs) -> list[tuple[str, dic last_punctuation != -1 and last_punctuation > self.CHARS_PER_TOKEN * min_tokens ): - chunk = chunk[: last_punctuation + 1] - i += len(tokenizer.encode_ordinary(chunk)) + kept = chunk[: last_punctuation + 1] + # Keep only whole source tokens whose decoded text stays within the + # punctuation cut, and advance the cursor by exactly that many + # tokens. The token straddling the cut (and everything after it) is + # then re-emitted in the next chunk instead of being skipped, which + # is what previously dropped characters. Token boundaries need not + # line up with the character cut, so we take the longest token prefix + # that is still a prefix of ``kept``. + consumed = 0 + for n in range(1, len(chunk_tokens) + 1): + if kept.startswith(tokenizer.decode(chunk_tokens[:n])): + consumed = n + else: + break + chunk = tokenizer.decode(chunk_tokens[:consumed]) if consumed else kept + # Guard against a chunk that would not advance the cursor, which would + # otherwise stall the loop forever. + i += max(consumed, 1) output.append((chunk, metadata)) return output diff --git a/python/pathway/xpacks/llm/tests/test_splitters.py b/python/pathway/xpacks/llm/tests/test_splitters.py index 88bf4b1de..706dfa23a 100644 --- a/python/pathway/xpacks/llm/tests/test_splitters.py +++ b/python/pathway/xpacks/llm/tests/test_splitters.py @@ -31,6 +31,20 @@ def test_tokencount(): assert_table_equality(result, input_table) +def test_tokencount_does_not_drop_characters(): + # Regression: when a chunk was cut at a punctuation mark, the cursor advanced + # by the re-encoded kept text, which could skip the tokens straddling the cut + # and silently drop characters. Concatenating the chunks must reconstruct the + # (unicode-normalized) input exactly. + import unicodedata + + splitter = TokenCountSplitter(min_tokens=1, max_tokens=3) + txt = "a.b.c.d.e.f.g.h.i.j.k.l." + chunks = [chunk for chunk, _ in splitter.chunk(txt)] + + assert "".join(chunks) == unicodedata.normalize("NFKC", txt) + + def test_recursive_from_encoding(): splitter = RecursiveSplitter( encoding_name="cl100k_base", chunk_size=30, chunk_overlap=0 From df560d3d41523299edb96b28362f7772cb9f65fa Mon Sep 17 00:00:00 2001 From: hylin Date: Tue, 6 Oct 2026 06:53:45 +0800 Subject: [PATCH 2/2] fix(xpacks/llm): locate token split boundaries by byte length --- CHANGELOG.md | 10 ++----- python/pathway/xpacks/llm/splitters.py | 28 ++++++++--------- .../xpacks/llm/tests/test_splitters.py | 30 +++++++++++++++++++ 3 files changed, 47 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 35f384e9c..facd8f7f5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,9 @@ This project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.htm ### Added - `pw.xpacks.llm.rerankers.LLMReranker` now accepts `call_kwargs`, the kwargs passed to each call of the LLM. The default is still `{"temperature": 0}`; pass `call_kwargs={}` to use the reranker with models that accept only the default `temperature`, such as Claude Opus 4.7 and newer (including Claude Opus 5.5) on Bedrock. +### Fixed +- `pathway.xpacks.llm.splitters.TokenCountSplitter` no longer skips source text when a punctuation cut falls inside a token. It locates the cut using token byte lengths and advances by exactly the whole source tokens emitted in the chunk. + ## [0.33.0] - 2026-09-18 ### Changed @@ -25,13 +28,6 @@ This project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.htm - The base installation no longer depends on `boto3`. Everything it backed works without it: the region of an S3 bucket is detected with a single anonymous request, and AWS credentials for the Delta Lake connector are resolved by the official AWS SDK chain built into the engine — environment variables, profile files (including SSO and assume-role), ECS and EC2 instance credentials. If your own code imports `boto3` and relied on Pathway pulling it in, add it to your project's dependencies. - Building a pipeline is now significantly faster: the location information attached to every expression and operator (used in the `Occurred here` part of error messages) is collected without materializing full stack traces. The speedup is most visible in code that constructs many small pipelines, such as unit test suites, where graph construction could previously dominate the run time — measured up to ~6x faster on pytest suites of small-data pipeline tests. Error messages are unchanged. -### Fixed -- `pw.io.rabbitmq.read` in the static mode no longer waits forever on a stream that holds no messages: such a read finishes with an empty table. The end of the stream is also waited for longer (15s instead of 5s) before the stream is believed empty, so a slow broker is less likely to turn a static read into an empty one. -- Output connectors with a `time` column (`pw.io.deltalake`, `pw.io.iceberg`, `pw.io.postgres`, `pw.io.mysql`, `pw.io.mssql`, `pw.io.sqlite`, `pw.io.duckdb`, `pw.io.mongodb`) no longer fail with a `worker panic: … TryFromIntError` (or an `IntOutOfRange` error) when they receive the rows that a `windowby`/`forget` buffer with a `delay` releases once the input ends. Such rows are written with `time` equal to the maximal 64-bit signed integer, marking the end of the stream. -- `pw.io.mqtt.read` with persistence enabled no longer loses messages on restart. With `qos` 1 or 2, a message is acknowledged to the broker only after a durable checkpoint covers it, and the broker session now survives restarts — so both the messages received shortly before a crash and the messages published while the pipeline was down are redelivered (at-least-once delivery). -- `pw.io.nats.read` with persistence enabled and a JetStream stream no longer loses the messages read between the last checkpoint and a crash: they are acknowledged only once a checkpoint covers them, so the server redelivers them after the restart (at-least-once delivery). -- `pathway.xpacks.llm.splitters.TokenCountSplitter` no longer drops characters when a chunk is cut at a punctuation mark. The cursor was advanced by the number of tokens the *kept* text re-encodes to, but token boundaries need not align with the character cut, so the token straddling the cut was skipped and the characters between the punctuation and the next token boundary were lost (e.g. `"a.b.c.d.e.f."` split with `max_tokens=3` dropped `f`). The splitter now advances by exactly the whole tokens kept in the chunk and re-emits the straddling token in the next chunk, so concatenating the chunks reconstructs the input. - ### Added - `pw.io.python.ConnectorSubject` now accepts `session_type="upsert"`: a row sent with an already-present primary key replaces the previous version instead of being an error, and the new public `delete` method removes a row. This is the intended mode for polling-style connectors that periodically re-send everything they see; it requires the schema to define a primary key. - `pw.io.postgres.write` and `pw.io.postgres.write_snapshot` now support PostGIS `geometry` and `geography` columns: a `str` column holding WKT (in any dimension, optionally with an `SRID=n;` prefix) or hex-encoded EWKB is converted to the binary EWKB form the server expects. `pw.io.postgres.read` reads such columns back as hex-encoded EWKB strings, so geometry values round-trip. diff --git a/python/pathway/xpacks/llm/splitters.py b/python/pathway/xpacks/llm/splitters.py index 3c754f1a5..3d94f70f9 100644 --- a/python/pathway/xpacks/llm/splitters.py +++ b/python/pathway/xpacks/llm/splitters.py @@ -267,21 +267,21 @@ def chunk(self, text: str, metadata: dict = {}, **kwargs) -> list[tuple[str, dic last_punctuation != -1 and last_punctuation > self.CHARS_PER_TOKEN * min_tokens ): - kept = chunk[: last_punctuation + 1] - # Keep only whole source tokens whose decoded text stays within the - # punctuation cut, and advance the cursor by exactly that many - # tokens. The token straddling the cut (and everything after it) is - # then re-emitted in the next chunk instead of being skipped, which - # is what previously dropped characters. Token boundaries need not - # line up with the character cut, so we take the longest token prefix - # that is still a prefix of ``kept``. - consumed = 0 - for n in range(1, len(chunk_tokens) + 1): - if kept.startswith(tokenizer.decode(chunk_tokens[:n])): - consumed = n - else: + cut_bytes = len(chunk[: last_punctuation + 1].encode("utf-8")) + token_bytes = 0 + cut_tokens = 0 + # Token boundaries can split a UTF-8 character. Count bytes rather + # than decoding every growing prefix to locate the punctuation cut. + for token in chunk_tokens: + token_bytes += len(tokenizer.decode_single_token_bytes(token)) + if token_bytes > cut_bytes: break - chunk = tokenizer.decode(chunk_tokens[:consumed]) if consumed else kept + cut_tokens += 1 + # If the first token straddles the cut, keep the full window so + # that the emitted text still matches the tokens we consume. + if cut_tokens: + consumed = cut_tokens + chunk = tokenizer.decode(chunk_tokens[:consumed]) # Guard against a chunk that would not advance the cursor, which would # otherwise stall the loop forever. i += max(consumed, 1) diff --git a/python/pathway/xpacks/llm/tests/test_splitters.py b/python/pathway/xpacks/llm/tests/test_splitters.py index 706dfa23a..afaa24ebc 100644 --- a/python/pathway/xpacks/llm/tests/test_splitters.py +++ b/python/pathway/xpacks/llm/tests/test_splitters.py @@ -3,6 +3,7 @@ from __future__ import annotations import pandas as pd +import pytest import pathway as pw from pathway.tests.utils import assert_table_equality @@ -45,6 +46,35 @@ def test_tokencount_does_not_drop_characters(): assert "".join(chunks) == unicodedata.normalize("NFKC", txt) +@pytest.mark.parametrize( + "txt", + [ + "Привет, мир. Это проверка разбиения текста! Работает ли оно? Да. " * 5, + "你好,世界。这是一个测试!它有效吗?是的. " * 10, + ], + ids=["russian", "chinese"], +) +def test_tokencount_does_not_duplicate_non_ascii(txt): + import unicodedata + + splitter = TokenCountSplitter() + chunks = [chunk for chunk, _ in splitter.chunk(txt)] + + assert "".join(chunks) == unicodedata.normalize("NFKC", txt) + + +def test_tokencount_preserves_token_straddling_first_punctuation_cut(): + # cl100k_base encodes "...)" as a single token. If the + # punctuation prefix has no whole token, emit the window intact. + splitter = TokenCountSplitter(min_tokens=0, max_tokens=1) + txt = "...) tail" + metadata = {"source": "example"} + chunks = splitter.chunk(txt, metadata) + + assert "".join(chunk for chunk, _ in chunks) == txt + assert all(chunk and meta == metadata for chunk, meta in chunks) + + def test_recursive_from_encoding(): splitter = RecursiveSplitter( encoding_name="cl100k_base", chunk_size=30, chunk_overlap=0