wellkilo opened a new issue, #1173:
URL: https://github.com/apache/flink-agents/issues/1173
### Search before asking
- [x] I searched the open issues and pull requests and found nothing similar.
### Description
`KafkaActionStateStore.put()` currently ignores the `Future<RecordMetadata>`
returned by `Producer.send()`:
```java
producer.send(kafkaRecord);
actionStates.put(stateKey, state);
producer.flush();
```
This catches synchronous `send()` or `flush()` failures, but not a
record-level failure reported asynchronously by Kafka after `send()` returns.
`flush()` waits for outstanding sends to finish; it does not surface the
exception stored in each send future.
As a result, a broker acknowledgement failure can be treated as a successful
durable-state write. The state is placed in the attempt-local cache and action
execution continues. A later checkpoint can capture Kafka end offsets that do
not include the failed record. After failover, recovery cannot find the
completed action state and may execute the action and its external side effects
again.
This affects both initial action-state writes and completed-action writes
because they share `KafkaActionStateStore.put()`. Java and Python agent paths
both use the same Java runtime store.
### How to reproduce
Kafka's `MockProducer` demonstrates the relevant API behavior
deterministically:
```java
MockProducer<String, String> producer =
new MockProducer<>(false, null, new StringSerializer(), new
StringSerializer());
Future<RecordMetadata> future =
producer.send(new ProducerRecord<>("topic", "key", "value"));
assertTrue(producer.errorNext(new RuntimeException("broker ack failed")));
assertDoesNotThrow(producer::flush);
assertThrows(ExecutionException.class, future::get);
```
A regression test against `KafkaActionStateStore.put()` can inject the same
failed future and assert that the method propagates the failure instead of
returning normally.
### Proposed fix
Wait for the future returned by the specific `send()` before updating the
in-memory cache. This keeps the existing synchronous persistence contract while
making record-level Kafka errors observable to `DurableExecutionManager`, which
already propagates persistence failures to fail the task.
The change does not require a public API, event schema, checkpoint format,
or configuration change.
### Version and environment
Current `main` (`2e97add7`, `0.4-SNAPSHOT`), Kafka client 4.0.0. The
behavior is not broker- or platform-specific.
### Are you willing to submit a PR?
- [x] I'm willing to submit a PR!
--
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]