sangkyoonnam commented on code in PR #1177:
URL: https://github.com/apache/flink-agents/pull/1177#discussion_r4235328733


##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -224,13 +224,16 @@ public void put(Object key, long seqNum, Action action, 
Event event, ActionState
         try {
             ProducerRecord<String, ActionState> kafkaRecord =
                     new ProducerRecord<>(topic, stateKey, state);
-            producer.send(kafkaRecord);
-            actionStates.put(stateKey, state);
+            producer.send(kafkaRecord).get();
             producer.flush();

Review Comment:
   For small records in otherwise idle batches, waiting on `get()` before 
`flush()` adds roughly the default linger delay. Puts are sequential, so the 
next put can't fill the batch while this one waits. kafka-clients 4.0 changed 
the `linger.ms` default from 0 to 5 ms, and `createProducerProp()` doesn't set 
it. The existing `flush()` is what bypassed linger. After `get()` returns, it 
no longer helps this record. `FlussActionStateStore` sets its writer batch 
timeout to zero for the same reason (L156-L158).
   
   I timed the three orderings with a plain `KafkaProducer` (kafka-clients 
4.0.0, `acks=all`, `retries=3`, string payloads) against a local 
`apache/kafka-native:3.8.0` container. 200 sequential sends per ordering, 
second of two rounds, ms per send iteration:
   
   ```
   main  send; flush         0.50
   PR    send.get; flush     6.88
   alt   send; flush; get    0.84
   ```
   
   A fresh action writes initial and completion state, and durable calls can 
add more synchronous writes, so the delay repeats per action. Flushing first 
bypasses linger, and `get()` still checks the record's result before the cache 
insert:
   
   ```java
   Future<RecordMetadata> ack = producer.send(kafkaRecord);
   producer.flush();
   ack.get();
   actionStates.put(stateKey, state);
   ```
   
   Setting `linger.ms` to 0 in `createProducerProp()` would also work and keeps 
your order. With the snippet, the broker-failure test still passes. The 
interrupt test fails only on `assertThat(flushCalled).isFalse()` at L277; 
flipping it to `isTrue()` keeps the ordering pinned. An interrupt while blocked 
in `flush()` produces Kafka's `InterruptException` wrapped in the generic 
`RuntimeException`, with the flag restored. That matches main today. Surfacing 
`InterruptedException` from that path would need a catch that translates 
Kafka's exception. I haven't tested a real interrupt. The flow and interaction 
table in the PR description would need the same update.



##########
runtime/API.md:
##########


Review Comment:
   Non-blocking. These two files repeat what the PR description already says, 
and nothing in the repo links to them. I'd remove them. If the contract needs a 
permanent home, the Javadoc on `KafkaActionStateStore.put` keeps it next to the 
code: for an initialized store, a normal return means the record was 
acknowledged and cached.



-- 
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