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](https://docs.confluent.io/cloud/current/ai/streaming-agents/agent-runtime-guide.html#context-management)
`tokens_management_strategy` (TRIM / SUMMARIZE past `max_tokens_threshold`,
off by default) and `max_iterations`; [Snowflake Cortex
Agents](https://docs.snowflake.com/en/user-guide/snowflake-cortex/cortex-agents-compact)
`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]