Skip to main content

Bootstrapping execution environment

Flinkboot automatically resolves, instantiates, and wires Apache Flink's StreamExecutionEnvironment directly from your YAML configuration with zero imperative plumbing.

Instead of handwriting hundreds of lines configuring state backends, checkpoints, exponential backoffs, and metrics, Flinkboot translates your declarative JobProperties into production-grade Flink execution settings in a single call.

// Instantiates and auto-configures the Flink environment from your YAML
StreamExecutionEnvironment env = boot.executionEnvironment(config.job());

1. How bootstrapping works​

When you invoke boot.executionEnvironment(JobProperties properties), Flinkboot applies the following lifecycle under the hood:

  1. Fail-Fast Invariant Checks: Verifies all cross-field constraints (e.g. exponential backoff boundaries, valid state backend classes, and checkpoint storage paths).
  2. Context Discovery (MiniCluster vs Cluster):
    • Local WebUI Mode (local-web-ui.enabled: true): Instantiates an embedded local MiniCluster with the Flink WebUI active (StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(...)).
    • Cluster Mode (local-web-ui.enabled: false or omitted): Calls Flink's native StreamExecutionEnvironment.getExecutionEnvironment(...) to automatically adapt to your runtime context (Standalone, Kubernetes, YARN, or local testing).
  3. Automated Setting Application: Injects RocksDB state backend, unaligned checkpointing, restart strategies, buffer timeouts, and watermark intervals directly onto Flink's native Configuration object.

2. Complete production YAML example​

Here is a complete YAML reference demonstrating all supported execution environment features:

# Configuration block mapped to JobProperties
name: "order-stream-processor"

environment:
execution:
runtime-mode: "STREAMING"
parallelism: 8
max-parallelism: 128
buffer-timeout: "PT0.05S"
auto-watermark-interval: "PT0.2S"
object-reuse: true

checkpointing:
enabled: true
interval: "PT30S"
mode: "EXACTLY_ONCE"
timeout: "PT2M"
min-pause-between-checkpoints: "PT5S"
max-concurrent-checkpoints: 1
externalized-checkpoint-cleanup: "RETAIN_ON_CANCELLATION"
unaligned-checkpoints: true
aligned-checkpoint-timeout: "PT1S"
storage-uri: "s3://flink-checkpoints/order-processor"

state-backend:
type: "ROCKSDB"
checkpoint-storage: "FILESYSTEM"
incremental: true
latency-tracking: true

restart-strategy:
type: "EXPONENTIAL_DELAY"
exponential-delay:
initial-backoff: "PT1S"
max-backoff: "PT1M"
backoff-multiplier: 2.0
reset-backoff-threshold: "PT1H"
jitter-factor: 0.1

savepoint-restore:
savepoint-path: "/mnt/savepoints/savepoint-0001"
allow-non-restored-state: false
restore-mode: "CLAIM"

local-web-ui:
enabled: false
port: 8081
bind-address: "localhost"

# Universal escape-hatch for arbitrary Flink Configuration options
properties:
taskmanager.memory.managed.fraction: "0.4"
pipeline.operator-chaining.enabled: "true"

3. Configuration reference specification​

Below is the complete specification of all properties supported by JobProperties.

A. Job metadata​

Property KeyTypeRequiredDescription
nameStringYesCanonical job name registered with Flink (PipelineOptions.NAME).
environmentObjectNoExecution environment settings (ExecutionEnvironmentProperties).

B. Execution options (environment.execution:)​

Property KeyTypeRequiredValidationDescription
runtime-modeEnumNoEnumExecution runtime mode: STREAMING, BATCH, or AUTOMATIC (ExecutionOptions.RUNTIME_MODE).
parallelismIntegerNo@PositiveDefault execution parallelism (CoreOptions.DEFAULT_PARALLELISM).
max-parallelismIntegerNo@PositiveMaximum parallelism for key groups rescale (PipelineOptions.MAX_PARALLELISM).
buffer-timeoutDurationNo@DurationMin(millis = 0)Buffer timeout (ExecutionOptions.BUFFER_TIMEOUT), e.g. "PT0.05S".
auto-watermark-intervalDurationNo@DurationMin(millis = 0)Periodic watermark emission interval (PipelineOptions.AUTO_WATERMARK_INTERVAL), e.g. "PT0.2S".
object-reuseBooleanNoBooleanEnable object reuse optimization (PipelineOptions.OBJECT_REUSE). Defaults to false.

C. Fault tolerance & checkpointing (environment.checkpointing:)​

Property KeyTypeRequiredValidationDescription
enabledBooleanNoBooleanMaster switch for checkpointing. Defaults to true when block is present.
intervalDurationNo@DurationMin(millis = 1)Time interval between checkpoints (CheckpointingOptions.CHECKPOINTING_INTERVAL), e.g. "PT30S".
modeEnumNoEnumConsistency mode: EXACTLY_ONCE or AT_LEAST_ONCE.
timeoutDurationNo@DurationMin(millis = 1)Maximum checkpoint duration before aborting (CheckpointingOptions.CHECKPOINTING_TIMEOUT).
min-pause-between-checkpointsDurationNo@DurationMin(millis = 0)Minimum rest duration between consecutive checkpoints.
max-concurrent-checkpointsIntegerNo@PositiveMaximum concurrent checkpoints allowed.
externalized-checkpoint-cleanupEnumNoEnumCleanup retention on cancel: RETAIN_ON_CANCELLATION, DELETE_ON_CANCELLATION, or NO_EXTERNALIZED_CHECKPOINTS.
unaligned-checkpointsBooleanNoBooleanEnable unaligned checkpoints (CheckpointingOptions.ENABLE_UNALIGNED).
aligned-checkpoint-timeoutDurationNo@DurationMin(millis = 0)Timeout before switching to unaligned checkpoints, e.g. "PT1S".
storage-uriStringNoStringTarget checkpoint storage URI, e.g. s3://bucket/checkpoints.

D. State backend & RocksDB (environment.state-backend:)​

Property KeyTypeRequiredValidationDescription
typeEnumNoEnumState backend type: ROCKSDB, HASHMAP, CHANGELOG, or CUSTOM.
checkpoint-storageEnumNoEnumStorage mechanism: JOBMANAGER or FILESYSTEM.
incrementalBooleanNoBooleanEnable incremental checkpoints for RocksDB (CheckpointingOptions.INCREMENTAL_CHECKPOINTS).
latency-trackingBooleanNoBooleanEnable latency tracking metrics for state access (StateBackendOptions.LATENCY_TRACK_ENABLED).
custom-classStringYes (if type == CUSTOM)StringFully qualified class name for custom state backend. Allowed only when type: CUSTOM.

E. Restart strategies (environment.restart-strategy:)​

The restart-strategy block accepts a type (NO_RESTART, FIXED_DELAY, FAILURE_RATE, EXPONENTIAL_DELAY, FALLBACK) and at most one matching sub-block.

Property KeyTypeRequiredValidationDescription
typeEnumNoEnumStrategy type: NO_RESTART, FIXED_DELAY, FAILURE_RATE, EXPONENTIAL_DELAY, or FALLBACK.
fixed-delayObjectNo@ValidParameters for FIXED_DELAY strategy.
failure-rateObjectNo@ValidParameters for FAILURE_RATE strategy.
exponential-delayObjectNo@ValidParameters for EXPONENTIAL_DELAY strategy.

1. Fixed delay (type: FIXED_DELAY)​

restart-strategy:
type: "FIXED_DELAY"
fixed-delay:
attempts: 3
delay: "PT10S"
Property KeyTypeRequiredValidationDescription
attemptsIntegerYes@PositiveOrZeroNumber of restart attempts before job failure.
delayDurationNo@DurationMin(millis = 0)Delay between restart attempts. Defaults to 0s.

2. Failure rate (type: FAILURE_RATE)​

restart-strategy:
type: "FAILURE_RATE"
failure-rate:
max-failures-per-interval: 5
failure-interval: "PT5M"
delay: "PT10S"
Property KeyTypeRequiredValidationDescription
max-failures-per-intervalIntegerYes@PositiveMaximum failures permitted within the time window.
failure-intervalDurationYes@DurationMin(millis = 1)Measurement window interval.
delayDurationNo@DurationMin(millis = 0)Delay between restart attempts. Defaults to 0s.

3. Exponential delay (type: EXPONENTIAL_DELAY)​

restart-strategy:
type: "EXPONENTIAL_DELAY"
exponential-delay:
initial-backoff: "PT1S"
max-backoff: "PT1M"
backoff-multiplier: 2.0
reset-backoff-threshold: "PT1H"
jitter-factor: 0.1
Property KeyTypeRequiredValidationDescription
initial-backoffDurationYes@DurationMin(millis = 1)Initial backoff duration.
max-backoffDurationYesGreater than or equal to initial-backoffMaximum cap for backoff duration.
backoff-multiplierDoubleNoStrictly > 1.0Exponential multiplier for consecutive failures. Defaults to 2.0.
reset-backoff-thresholdDurationNo@DurationMin(millis = 1)Failure-free uptime required to reset backoff. Defaults to PT1H.
jitter-factorDoubleNo0.0 <= x <= 1.0Random jitter ratio added to backoff. Defaults to 0.1.

F. Savepoint recovery (environment.savepoint-restore:)​

savepoint-restore:
savepoint-path: "/mnt/savepoints/savepoint-0001"
allow-non-restored-state: false
restore-mode: "CLAIM"
Property KeyTypeRequiredValidationDescription
savepoint-pathStringYes@NotBlankPath to savepoint or initial checkpoint directory.
allow-non-restored-stateBooleanNoBooleanStart even if savepoint contains unmapped subtask state.
restore-modeEnumNoEnumRestore mode: CLAIM, NO_CLAIM, or LEGACY.

G. Local dev WebUI (environment.local-web-ui:)​

local-web-ui:
enabled: true
port: 8081
bind-address: "localhost"
Property KeyTypeRequiredValidationDescription
enabledBooleanYes (if block present)@NotNullStarts a local Flink MiniCluster with the WebUI dashboard active during IDE testing.
portIntegerNo@Range(min = 0, max = 65535)WebUI REST port. Defaults to 8081. Set to 0 for dynamic port allocation.
bind-addressStringNo@NotBlankWebUI host bind address. Defaults to localhost.

Enabling local-web-ui.enabled: true requires org.apache.flink:flink-runtime-web on the classpath:

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-runtime-web</artifactId>
<scope>provided</scope>
</dependency>

When running directly from an IDE (such as IntelliJ IDEA), ensure that the option "Include dependencies with 'Provided' scope" is enabled in your run configuration so the WebUI classes are present on the execution classpath.

If set to true inside a remote Flink cluster, Flinkboot fails fast with an UnsupportedExecutionEnvironmentException.


H. Universal escape hatch (environment.properties:)​

Arbitrary Flink configuration key-value pairs applied directly onto Flink's native Configuration object:

environment:
properties:
taskmanager.memory.managed.fraction: "0.4"
pipeline.operator-chaining.enabled: "true"
Property KeyTypeRequiredDescription
propertiesMap<String, String>NoFree-form key-value map mapped directly to Flink's native Configuration. Properties defined here take direct precedence over typed YAML properties in case of conflict.