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]