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]

Reply via email to