Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,17 @@ public enum ConditionEvaluationFailureStrategy {
public static final ConfigOption<Boolean> KAFKA_ACTION_STATE_TOMBSTONE_ENABLED =
new ConfigOption<>("kafkaActionStateTombstoneEnabled", Boolean.class, false);

/**
* The separate, single-partition Kafka topic that stores committed checkpoint-aligned cleanup
* boundaries. It must use {@code cleanup.policy=compact} without delete retention. Setting this
* option enables boundary enforcement during recovery. It must not be the action-state data
* topic, and it must be dedicated to one job's recovery history. The action-state data topic
* must use {@code cleanup.policy=compact,delete}, {@code retention.ms=-1}, and {@code
* retention.bytes=-1} so only reviewed cleanup plans advance its prefix.
*/
public static final ConfigOption<String> KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC =
new ConfigOption<>("kafkaActionStateCleanupControlTopic", String.class, null);

/** The config parameter specifies the Fluss bootstrap servers. */
public static final ConfigOption<String> FLUSS_BOOTSTRAP_SERVERS =
new ConfigOption<>("flussBootstrapServers", String.class, "localhost:9123");
Expand Down
24 changes: 24 additions & 0 deletions docs/content/docs/operations/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -183,6 +183,30 @@ Here are the configuration options for Kafka-based Action State Store.
| `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 on a compacted action-state topic. Off by default: pruning does not invalidate older restore points, but the topic continues to grow. When enabled, the checkpoint whose completion triggers pruning remains usable, but restoring an earlier checkpoint or savepoint may replay later tombstones and re-execute already completed actions. Enable only if the job never restores from earlier checkpoints or savepoints, or if re-executing actions is acceptable. |
| `kafkaActionStateCleanupControlTopic` | (none) | String | Separate, single-partition topic containing committed checkpoint-aligned cleanup boundaries. It must use `cleanup.policy=compact` without delete retention, differ from the action-state topic, and be dedicated to the same job recovery history. Setting it enables boundary enforcement during recovery. It cannot be combined with `kafkaActionStateTombstoneEnabled`. |

##### Checkpoint-aligned Kafka cleanup

Checkpoint-aligned cleanup is an explicit administrative operation. `KafkaActionStateCleanupTool` provides the same command-line workflow for Java and Python jobs. Run it with the Flink Agents runtime JAR and the matching Flink State Processor API JAR on the classpath:

```text
plan --checkpoint PATH (--operator-uid UID | --operator-uid-hash HASH) --output FILE
apply --plan FILE --bootstrap-servers SERVERS --control-topic TOPIC [--replication-factor N]
```

The `plan` command reads all recovery markers from the selected checkpoint or savepoint through Flink's State Processor API, creates a deterministic plan using the earliest required offset per partition, and refuses to overwrite an existing output file. Legacy map markers can still be restored by a job, but cannot authorize deletion because they do not identify the physical Kafka topic. The `apply` command verifies the content-derived plan ID before contacting Kafka. Plan files must contain exactly one JSON document; trailing content is rejected. The optional `--replication-factor` controls replication when creating the control topic, defaults to `1`, and must be between `1` and `32767`.

Review and retain the plan's deterministic JSON before applying it. The coordinator verifies that the selected offsets are still available, writes `COMMITTED` before calling Kafka `deleteRecords`, verifies every resulting beginning offset, and then writes `APPLIED`. Reapplying the same plan retries an interrupted committed operation without changing its boundary. Set `kafkaActionStateCleanupControlTopic` on the job only after the first plan is committed; recovery requires the configured topic to exist and contain a committed boundary, and fails closed if the topic is missing or empty.

Physical prefix cleanup requires the action-state data topic to use `cleanup.policy=compact,delete`, `retention.ms=-1`, and `retention.bytes=-1`. The store uses those settings when it creates a new topic. Before applying cleanup to an existing topic, alter it to those values. The `delete` policy enables Kafka's `deleteRecords` API, while both retention limits remain disabled so Kafka cannot independently retire records that a supported checkpoint still needs. Apply validates these effective settings before writing `COMMITTED` and rechecks them immediately before and after deletion.

For the first cleanup on an existing job, stop the job, create and apply the plan from the recovery point that will become the oldest supported one, and restart from that same point with `kafkaActionStateCleanupControlTopic` configured. This avoids a failover window in which the old running job does not yet know the committed boundary. Once the job is running with boundary enforcement enabled, later forward-only plans may be applied while it runs. Run only one `apply` operation at a time; concurrent incomparable plans fail closed and require operator intervention. Do not recreate the action-state topic or change its partitions while an apply operation is running: the coordinator checks the topic ID and partition set before and after deletion, but Kafka addresses `deleteRecords` by topic name and cannot atomically fence topic lifecycle changes.

The action-state topic and control topic must be dedicated to one job's recovery history. Once any boundary has been committed, every future run and restore of that history must retain the same control-topic configuration; omitting or changing it removes logical boundary enforcement. The boundary can move only forward. A checkpoint whose marker is below the committed boundary is rejected even if Kafka has not finished physical deletion.

Per-key tombstones and checkpoint-aligned cleanup are mutually exclusive. Tombstones are not tied to the selected recovery boundary and could invalidate a checkpoint that checkpoint-aligned cleanup promises to retain.

To migrate a job that previously emitted tombstones, first stop it cleanly so its Kafka producer closes and all accepted tombstone sends finish. Restart from the newest recovery point that is valid under the tombstone mode's existing recovery trade-off, with `kafkaActionStateTombstoneEnabled=false` and no cleanup control topic configured. After that attempt completes a new checkpoint or savepoint, stop it again, create and apply the first cleanup plan from that new recovery point, and restart from the same point with `kafkaActionStateCleanupControlTopic` configured. The new marker is after every old tombstone; selecting an older recovery point cannot provide the same guarantee because replay may still encounter those tombstones.

#### Fluss-based Action State Store

Expand Down
4 changes: 4 additions & 0 deletions docs/content/docs/operations/deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,10 @@ See [Action State Store Configuration]({{< ref "docs/operations/configuration#ac
**Note**: Enabling Kafka action-state tombstones can invalidate checkpoints or savepoints older than the prune and cause completed actions to execute again. See [Action State Store Configuration]({{< ref "docs/operations/configuration#action-state-store" >}}) for the recovery trade-off.
{{< /hint >}}

{{< hint warning >}}
**Note**: A committed checkpoint-aligned Kafka cleanup boundary permanently makes older recovery points unsupported. Every subsequent run and restore must use the same cleanup control topic so the runtime can reject those recovery points before replay. See [Checkpoint-aligned Kafka cleanup]({{< ref "docs/operations/configuration#checkpoint-aligned-kafka-cleanup" >}}).
{{< /hint >}}

{{< hint info >}}
**Note**: Exactly-once action consistency is guaranteed only if, after recovering from the same checkpoint, inputs for each key arrive in the same order as before recovery. If this ordering requirement is not met, the system falls back to exactly-once output consistency.
{{< /hint >}}
6 changes: 6 additions & 0 deletions python/flink_agents/api/core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,12 @@ class AgentConfigOptions:
default=False,
)

KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC = ConfigOption(
key="kafkaActionStateCleanupControlTopic",
config_type=str,
default=None,
)

FLUSS_BOOTSTRAP_SERVERS = ConfigOption(
key="flussBootstrapServers",
config_type=str,
Expand Down
3 changes: 3 additions & 0 deletions python/flink_agents/api/tests/test_core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,9 @@ def test_agent_config_options_are_explicitly_declared() -> None:
options = _collect_config_options(AgentConfigOptions)
assert options["BASE_LOG_DIR"].get_key() == "baseLogDir"
assert options["KAFKA_BOOTSTRAP_SERVERS"].get_default_value() == "localhost:9092"
cleanup_control_topic = options["KAFKA_ACTION_STATE_CLEANUP_CONTROL_TOPIC"]
assert cleanup_control_topic.get_key() == "kafkaActionStateCleanupControlTopic"
assert cleanup_control_topic.get_default_value() is None
assert options["EVENT_LOG_LEVEL"].get_default_value() is EventLogLevel.STANDARD
assert options["EVENT_LOG_TRACE_ENABLED"].get_default_value() is False
condition_failure = options["CONDITION_EVALUATION_FAILURE_STRATEGY"]
Expand Down
20 changes: 19 additions & 1 deletion runtime/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,12 @@ under the License.
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-state-processor-api</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java</artifactId>
Expand Down Expand Up @@ -122,6 +128,18 @@ under the License.
<artifactId>kafka-clients</artifactId>
<version>${kafka.version}</version>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>kafka</artifactId>
<version>1.21.4</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>commons-codec</groupId>
<artifactId>commons-codec</artifactId>
<version>1.16.1</version>
<scope>test</scope>
</dependency>
<!-- fluss client -->
<dependency>
<groupId>org.apache.fluss</groupId>
Expand Down Expand Up @@ -261,4 +279,4 @@ under the License.
</build>
</profile>
</profiles>
</project>
</project>
Loading