wellkilo opened a new pull request, #1177:
URL: https://github.com/apache/flink-agents/pull/1177
<!-- Please add the relevant components in the PR title. -->
Linked issue: #1173
### Purpose of change
A caller no longer receives a successful return from
`KafkaActionStateStore.put(...)` when Kafka
later rejects the state record. The asynchronous broker failure now fails
the task before the
attempt-local cache changes, preventing a checkpoint from treating an
action-state update as
durable when its Kafka record was never acknowledged.
Runtime flow:
1. `put(...)` encodes the composite state key and sends the `ProducerRecord`.
2. It waits on the `Future<RecordMetadata>` returned by that specific
`send()` call.
3. After the acknowledgement succeeds, it flushes the producer and updates
the attempt-local
cache.
4. Only then does `put(...)` return successfully.
#### Key decisions
- Wait on the send future instead of relying on `flush()`. Kafka's
`MockProducer` demonstrates
that `flush()` can return normally while the per-record future retains an
asynchronous failure.
- Keep the existing synchronous persistence contract: successful return
still means both the
acknowledgement and flush completed.
- Re-throw `InterruptedException` with the thread's interrupt flag restored,
rather than wrapping
it as a generic Kafka persistence failure.
- Do not add a new acknowledgement timeout or configuration. This patch only
makes the existing
wait observable; it does not change the configured producer behavior.
#### Interaction decisions
| Send / acknowledgement / flush outcome | Result |
| --- | --- |
| `send()` throws synchronously | `RuntimeException`; cache unchanged |
| Send future completes exceptionally | `RuntimeException` with the broker
failure as root cause; flush and cache update are skipped |
| Acknowledgement succeeds, `flush()` fails | `RuntimeException`; cache
update is skipped |
| Acknowledgement wait is interrupted | `InterruptedException`; interrupt
flag restored; flush and cache update are skipped |
| Acknowledgement and flush succeed | Cache is updated, then `put(...)`
returns normally |
#### Behavioral contracts
- A successful `put(...)` means the specific record was acknowledged, the
producer was flushed,
and the attempt-local cache contains the state.
- An asynchronous acknowledgement failure is propagated to the caller and
does not populate the
attempt-local cache.
- An interrupted acknowledgement wait preserves
`Thread.currentThread().isInterrupted()` and
throws `InterruptedException`.
- The method does not update the cache on send, acknowledgement, or flush
failure.
#### Failure behavior
Kafka send and flush exceptions are surfaced as
`RuntimeException("Failed to send action state to Kafka", cause)`.
Interruption is not wrapped:
it is re-thrown with the interrupt status preserved. A null producer retains
the pre-existing
behavior of logging and returning.
### Tests
| Behavioral contract | Evidence |
| --- | --- |
| Asynchronous broker failure fails `put(...)` and leaves the cache
unchanged | `testPutPropagatesAsyncBrokerFailureBeforeCachingState` |
| Interruption is propagated, the interrupt flag is restored, and the cache
is unchanged | `testPutPreservesInterruptWhileWaitingForBrokerAcknowledgement` |
| The successful send, flush, and cache path remains intact |
`testPutActionState` and the existing `KafkaActionStateStoreTest` suite |
Local verification completed on the branch: `KafkaActionStateStoreTest`
(51/51), repository
lint/Spotless, and Apache RAT all pass.
Not verified: live-broker failover/recovery and a never-completing broker
acknowledgement. The
regression uses `MockProducer` to make the record-level failure
deterministic; no acknowledgement
timeout is introduced by this change.
<details>
<summary>Implementation invariants and supporting evidence</summary>
- `actionStates.put(...)` is reached only after `future.get()` and
`producer.flush()` complete
successfully.
- The `InterruptedException` catch restores the flag before re-throwing the
same exception.
- The existing `ActionStateStore.put(...)` signature already declares
`throws Exception`, so no
public interface change is required.
- The Kafka record schema and the attempt-local cache key remain unchanged.
</details>
### API
No public API signature, event schema, checkpoint format, or configuration
change is introduced.
The caller-visible behavior change is limited to propagating failures that
were previously hidden
by ignoring the send future.
### Documentation
- [ ] `doc-needed`
- [ ] `doc-not-needed`
- [x] `doc-included`
### Was this patch authored or co-authored using generative AI tooling?
- [x] Yes
- [ ] No
Generated-by: TraeCode (GPT-5)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]