GitHub user purushah edited a discussion: [Discussion] Context compaction for 
Flink Agents

Hi all. I'd like to propose adding **context compaction** to Flink Agents and 
agree on its shape here before opening PRs. Below: the failure, what we found 
when we measured the obvious fixes, the strategies we propose, and our 
questions.

## The failure

A ReAct-style agent running as a Flink Agents job (PyFlink 2.3.0, 
`ActionExecutionOperator`, the built-in chat model action, 
`OpenAIChatModelSetup`, gpt-4o-mini). One input event, one tool returning 40 
rows (~2.6k tokens). To reach the limit quickly the setup forces a tool call on 
every round (`tool_choice="required"`, sequential calls), so this loop cannot 
end on its own; the un-forced agent in the second bullet below hits the same 
wall. Every round the built-in loop re-sends every earlier tool result. After 
48 rounds and 53 seconds:

```
2026-09-28 16:27:38,993 WARN  org.apache.flink.runtime.taskmanager.Task
  - OverflowReActAgent -> Map, Map -> Sink: Print to Std. Out (1/1)#0 switched 
from RUNNING to FAILED

java.lang.RuntimeException: 
org.apache.flink.agents.runtime.operator.ActionExecutionOperator$ActionTaskExecutionException:
 Failed to execute action task
Caused by: 
org.apache.flink.agents.runtime.operator.ActionExecutionOperator$ActionTaskExecutionException:
 Failed to execute action task
Caused by: pemja.core.PythonException: <class 'openai.BadRequestError'>: Error 
code: 400 - {'error': {'message': "This model's maximum context length is 
128000 tokens. However, you requested 129323 tokens (127719 in the messages, 
104 in the functions, and 1500 in the completion). Please reduce the length of 
the messages, functions, or completion.", 'type': 'invalid_request_error', 
'param': 'messages', 'code': 'context_length_exceeded'}}

2026-09-28 16:27:39,079 INFO  
org.apache.flink.runtime.executiongraph.ExecutionGraph
  - Job react-overflow-demo (e724a3bcc0350dffc778859b4cf97635) switched from 
state RUNNING to FAILING.
org.apache.flink.runtime.JobException: Recovery is suppressed by 
NoRestartBackoffTimeStrategy
```

```
JOB FAILED after 53s
tool invocations executed inside the operator: 48
```

- **It repeats.** Since #1149 the error becomes a failed `ChatResponseEvent`; 
retrying an unchanged oversized request cannot fix it, and depending on the 
durable-execution configuration recovery either replays the recorded failure or 
repeats the work (the run above disables restarts). Prompt caching does not 
help, cached tokens still count.
- **It happens unforced.** A data-driven agent (list accounts, list 
transactions, fetch every flagged one, submit a report) hit the same 400 on its 
40th call at 131,980 tokens.
- **It happens without tools.** An agent that appends each event and reply to a 
short-term-memory list and sends the list every call was rejected on the 61st 
event of one key at 129,968 tokens. No tool in the job.
- **Growth hurts first.** Unbounded, the tool-loop agent abandoned its plan 
after 14 of 60 lookups near 39k tokens; the history agent answered a direct 
question with a verdict at 97k tokens.

## Why there is no bound today

- **Inside one run**, the chat action's tool loop re-sends every tool result 
every round: no round cap, no budget, no truncation, no clearing.
- **Across runs**, there is no framework-owned conversation. Agents manage 
cross-event history themselves in short-term memory; `MemoryObject.set` can 
replace stored history (there is no `remove`, see the TODO in 
`chat_model_action.py`), but the framework provides no conversation-compaction 
policy.
- **Sub-agents as tools** (#1114, #1132) add one result per delegation to the 
parent's loop.

The existing controls (event-logger payload truncation, the routing judge's 
budget, and since #1040 detection of truncated *output*) do not impose a total 
input-context budget on the ordinary chat loop.

## What the cheap fixes preserve

We measured the cheap fixes before proposing anything, on two synthetic 
workloads with known answers and on the data-driven agent above; numbers on 
request. Full history gives wrong answers while it still fits, then fails hard. 
In these experiments, masking old tool results reduced prompt growth but could 
remove information needed later. Summarisation achieved higher answer recall 
than masking, with additional model calls and changes in task behaviour. Inside 
Flink, masking, per-result truncation and a round cap took the data-driven 
agent to 41 rounds with no provider error, but it never produced an accepted 
report: avoiding overflow does not guarantee task completion.


## State instead of a transcript

We also ran the 70-event key from the failure above with a small working state 
per key instead of a transcript, updated by the model on each event and sent in 
place of history. The request stays at one event plus one state object, about 
2.4k tokens, for all 70 events, against a 400 on event 61 for the transcript. 
The risk moves to state quality: a free-form state kept the BR whitelist but 
dropped its ticket id, which cannot be recovered from the state alone; 
prompting with an explicit schema kept both complete fact and ticket pairs; 
with both facts correctly stored, gpt-4o-mini passed two of four probes (one 
answer was incorrect, another omitted the required ticket) and gpt-4o passed 
all four in this single run. These are illustrative runs; the state arms used 
additional instructions and a larger output budget than the transcript arm. A 
working state gives no verbatim recall, so tasks that need exact historical 
details need retained originals and a retrieval path.


## Proposal: nine strategies, four groups, all off by default

Masking and truncation change only the model-visible messages and preserve the 
active request's stored tool-call payloads and tool_call_id pairing; the event 
log keeps what its truncation setting allows. The round cap changes control 
flow, group B changes stored state, and native compaction adds continuation 
state.

**A — deterministic trimming, no model call** (a prototype exists in Java and 
Python with tests and docs; it needs a hardening pass before a PR)

- **1. Round-based result masking**, `chat.tool-results.keep-last-rounds`: keep 
the last N rounds' results verbatim, replace older ones with a one-line 
placeholder (tool, arguments, omitted size). Assistant tool-call messages and 
placeholders stay, so the prompt still grows slowly (about 53 tokens per round 
in our run); this trims growth rather than bounding it. Peers: LangChain 
`ClearToolUsesEdit`, Anthropic `clear_tool_uses`, Microsoft 
`ToolResultCompactionStrategy`.
- **2. Per-result truncation**, `chat.tool-results.max-chars`: trim one 
oversized result in the middle. Precedent: core Flink `context-overflow-action`.
- **3. Round cap**, `chat.max-tool-rounds`: fail through #1149's path after N 
rounds without a final answer. Peers: OpenAI Agents SDK `max_turns`, Confluent 
`max_iterations`.

**B — state instead of transcript, no separate compaction call**

- **4. Bounded working state per key**: `MemoryObject.remove` (`set` already 
overwrites), an append primitive, and a documented pattern with the example 
above (typed object per key, schema validation, hard size cap). Stored in Flink 
managed keyed state and included in configured checkpoint recovery (not 
exercised in the run above). Venue: #1055, #1086, #1056.

**C — a model compresses**

- **5. Budget-triggered rolling summary** with a configured summariser (a 
routable chat-model resource), last K turns verbatim, sliding window as 
fallback, folded in batches to reduce how often the cached prefix changes 
(compaction invalidates cache reuse from the changed prefix onward).
- **6. Provider-native compaction** (Anthropic compaction blocks, OpenAI 
Responses compaction), opt-in per connection; needs provider-specific connector 
support (our OpenAI connection uses Chat Completions today, so this needs 
Responses support) and stores the complete continuation representation, for 
OpenAI the returned canonical window including retained items, in keyed state.

**D — a model decides (direction only, built on the routing judge)**

- **7. Judged retention**: a cheap judge marks old results as needed or 
droppable.
- **8. Judged summarisation target**: a judge lists what must survive verbatim 
before summarising.
- **9. Compaction routing**: a judge picks the rung (nothing, mask, summarise, 
stop).

For 5–9: durably recorded summariser and judge outcomes can be reused on 
replay; calls without a recoverable recorded outcome may execute again. The 
current action API has no background-compaction lifecycle, so adding one needs 
explicit coordination with per-key state and recovery. Not proposed: sliding 
window as primary, retrieval as compaction, token pruning or KV-cache methods; 
code-mode tools and sub-agents are patterns to document, not options.

## Why not just …

- *return less from tools?* Always, but we cannot assume third-party tools 
return concise results; the framework can still limit what reaches the model.
- *cap cost?* #1065 proposes per-agent token and cost budgets; 
`max-tool-rounds` is its loop-level sibling.
- *let the provider hold the conversation?* Does not remove context costs and 
needs explicit coordination with checkpoint recovery.
- *use a larger-context model?* Headroom, not a bound.
- *use sub-agents?* Right for long tasks (#1114, #1132, #1138); each loop still 
needs a bound.

## Questions

1. Is `chat.tool-results.keep-last-rounds` the right name and level (execution 
option vs chat-model setup field)?
2. Should the placeholder text be configurable in v1?
3. Should `max-tool-rounds` offer a `finish` mode (one last call without tools) 
besides failing?
4. Is `remove` + `append` on `MemoryObject` with a documented pattern enough 
for group B in 0.4, or should the framework offer a typed working-state helper?
5. Should the summariser default to the agent's own chat model or require an 
explicit resource?
6. Any objection to opening the group A PR after the hardening pass, with B–D 
as follow-up issues linked from this discussion?

## References

- Related: #1065, #1149, #1114, #1132, #1138, #1139, #1055, #1086, #1056, #876, 
#1146, #858.
- Prior art: The Complexity Trap (arXiv 2508.21433); OpenHands condensers; 
LangChain context editing and summarization middleware; Anthropic context 
editing and compaction; OpenAI Responses compaction; Google ADK 
`EventsCompactionConfig`; Microsoft Agent Framework compaction; Confluent 
Streaming Agents `tokens_management_strategy` (TRIM / SUMMARIZE past 
`max_tokens_threshold`, off by default) and `max_iterations`; Snowflake Cortex 
Agents `agent:compact` (preview); core Flink `max-context-size` / 
`context-overflow-action`.
- Measurements, traces, redacted TaskManager logs and the exact build and 
configuration for each run will be attached to the group A PR; the 
working-state agents to the group B issue.


GitHub link: https://github.com/apache/flink-agents/discussions/1202

----
This is an automatically sent email for [email protected].
To unsubscribe, please send an email to: [email protected]

Reply via email to