rob-9 commented on code in PR #885:
URL: https://github.com/apache/flink-agents/pull/885#discussion_r3834121683
##########
runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java:
##########
@@ -276,30 +282,64 @@ public void rebuildState(List<Object> recoveryMarkers) {
public void pruneState(Object key, long seqNum) {
LOG.debug("Pruning state for key: {} up to sequence number: {}", key,
seqNum);
- // Remove states from in-memory cache for this key up to the specified
sequence
- // number
- actionStates
- .entrySet()
- .removeIf(
- entry -> {
- String stateKey = entry.getKey();
- // Extract key and sequence number from the state
key
- // State key format: "key_seqNum_action_event"
- if (stateKey.startsWith(key.toString() + "_")) {
- try {
- List<String> parts =
ActionStateUtil.parseKey(stateKey);
- if (parts.size() >= 2) {
- long stateSeqNum =
Long.parseLong(parts.get(1));
- return stateSeqNum <= seqNum;
- }
- } catch (NumberFormatException e) {
+ // Collect state keys belonging to this key with sequence number <=
seqNum. The parsed
+ // key part must match exactly: prefix matching alone would let
pruning key "a_1" match
+ // state keys of the distinct key "a" (whose keys also start with
"a_1_").
+ String keyStr = key.toString();
+ String keyPrefix = keyStr + "_";
+ List<String> keysToPrune = new ArrayList<>();
+ for (String stateKey : actionStates.keySet()) {
+ if (!stateKey.startsWith(keyPrefix)) {
+ continue;
+ }
+ try {
+ List<String> parts = ActionStateUtil.parseKey(stateKey);
+ if (parts.get(0).equals(keyStr) &&
Long.parseLong(parts.get(1)) <= seqNum) {
+ keysToPrune.add(stateKey);
+ }
+ } catch (IllegalArgumentException e) {
+ LOG.warn(
+ "Cannot parse state key: {}. The entry cannot be
pruned and will be "
+ + "retained in memory and in the topic.",
+ stateKey,
+ e);
+ }
+ }
+
+ // Send tombstones to Kafka so log compaction can reclaim storage;
opt-in because
+ // tombstones break replay when restoring a checkpoint/savepoint older
than the prune
+ // (see KAFKA_ACTION_STATE_TOMBSTONE_ENABLED). Send failures surface
asynchronously,
+ // so report them via callback; the records then persist until manual
cleanup.
+ if (tombstoneEnabled && producer != null && !keysToPrune.isEmpty()) {
+ try {
+ for (String stateKey : keysToPrune) {
+ producer.send(
+ new ProducerRecord<>(topic, stateKey, null),
+ (metadata, exception) -> {
+ if (exception != null) {
LOG.warn(
- "Failed to parse sequence number
from state key: {}",
- stateKey);
+ "Failed to send tombstone record
for state key: {}. "
+ + "The record will persist
in the topic "
+ + "until manual cleanup.",
+ stateKey,
+ exception);
}
- }
- return false;
- });
+ });
+ }
+ producer.flush();
Review Comment:
removed, Kafka ordering is preserved by the state key, and delivery failures
remain reported through the callback.
##########
api/src/main/java/org/apache/flink/agents/api/configuration/AgentConfigOptions.java:
##########
@@ -64,6 +64,19 @@ public class AgentConfigOptions {
public static final ConfigOption<Integer>
KAFKA_ACTION_STATE_TOPIC_REPLICATION_FACTOR =
new ConfigOption<>("kafkaActionStateTopicReplicationFactor",
Integer.class, 1);
+ /**
+ * The config parameter determines whether pruning sends tombstone
(null-valued) records to the
+ * Kafka action state topic so log compaction can reclaim pruned keys.
Defaults to {@code
+ * false}: without tombstones the topic grows unboundedly, but restoring
any checkpoint or
+ * savepoint replays correctly. When enabled, restoring from the latest
completed checkpoint is
+ * unaffected, but restoring an older checkpoint or savepoint may replay
tombstones written
+ * after that restore point, erasing action state the replay still needs
and causing already
+ * completed actions to re-execute. Enable only if the job never restores
from non-latest
+ * checkpoints or savepoints, or if re-executing actions is acceptable.
+ */
+ public static final ConfigOption<Boolean>
KAFKA_ACTION_STATE_TOMBSTONE_ENABLED =
+ new ConfigOption<>("kafkaActionStateTombstoneEnabled",
Boolean.class, false);
Review Comment:
agreed leaving this PR scoped to opt-in cleanup; checkpoint-aligned auto
cleanup requires a separate design, will follow-up with an issue!
##########
docs/content/docs/operations/configuration.md:
##########
@@ -168,6 +168,7 @@ Here are the configuration options for Kafka-based Action
State Store.
| `kafkaActionStateTopic` | (none) | String |
The config parameter specifies the Kafka topic for action state. |
| `kafkaActionStateTopicNumPartitions`| 64 | Integer |
The config parameter specifies the number of partitions for the Kafka action
state topic. |
| `kafkaActionStateTopicReplicationFactor` | 1 | Integer |
The config parameter specifies the replication factor for the Kafka action
state topic. |
+| `kafkaActionStateTombstoneEnabled` | false | Boolean |
Whether pruning sends tombstone records so log compaction can reclaim pruned
keys. Off by default: without tombstones the topic grows unboundedly, but
restoring any checkpoint or savepoint replays correctly. When enabled,
restoring from the latest completed checkpoint is unaffected, but restoring an
older checkpoint or savepoint may replay tombstones written after that restore
point and re-execute already completed actions. Enable only if the job never
restores from non-latest checkpoints or savepoints, or if re-executing actions
is acceptable. |
Review Comment:
updated.
##########
runtime/src/test/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStoreTest.java:
##########
@@ -196,6 +207,129 @@ void testPruneState() throws Exception {
assertNull(
actionStates.get(ActionStateUtil.generateKey(TEST_KEY, 2L,
testAction, testEvent)));
assertNotNull(actionStateStore.get(TEST_KEY, 3L, testAction,
testEvent));
+
+ // Assert - tombstones should have been sent to Kafka
+ var history = mockProducer.history();
+ assertThat(history).hasSize(2);
Review Comment:
addressed. `testPruneState` now covers cache eviction, while the dedicated
tests own tombstone and default-off assertions.
--
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]