# Flinkboot Complete Documentation (Version 0.5.0)
> Bootstrapping & Reliability Framework for Apache Flink.
> Version: 0.5.0
> https://flinkboot.com | GitHub: https://github.com/Sekelenao/Flinkboot
================================================================================
DOCUMENT: Overview (/docs/)
================================================================================
import Tabs from '@theme/Tabs';
import TabItem from '@theme/TabItem';
# Overview
Building Apache Flink applications should feel as clean, safe, and productive as writing modern Spring Boot backend services. Once you experience declarative bootstrapping, fail-fast configuration, and zero-Kryo guarantees, you will never write manual Flink boilerplate again.
---
## 1. Type-safe configuration DTOs
Flinkboot ships with pre-built, production-tested DTOs for the execution environment (`JobProperties`) and official connectors (`KafkaSourceProperties`, `FlussSourceProperties`). Simply assemble them with your own domain records or POJOs with full Jakarta Bean Validation (`@NotNull`, `@Valid`), with no manual Jackson parsing required.
```java
// Composed application configuration record
public record AppConfig(
// Built-in Flinkboot DTO: covers RocksDB, restart strategies, metrics, checkpoints
@Valid @NotNull @JsonProperty("job") JobProperties job,
// Built-in Flinkboot DTO: covers brokers, topics, offset strategies, vendor escape-hatch
@Valid @NotNull @JsonProperty("kafka-source") KafkaSourceProperties kafkaSource,
// Your custom business domain settings (records or POJOs with Bean Validation)
@Valid @NotNull @JsonProperty("alerting") AlertingProperties alerting
) {}
```
```java
// Untyped, error-prone manual Jackson tree traversal
ObjectMapper mapper = new ObjectMapper(new YAMLFactory());
JsonNode root = mapper.readTree(new File("application.yaml"));
// Cryptic NullPointerExceptions at runtime if a single key is missing or misspelled
String jobName = root.path("job").path("name").asText();
int parallelism = root.path("job").path("parallelism").asInt(1);
long checkpointInterval = root.path("job").path("checkpointing").path("interval").asLong();
String brokers = root.path("kafka-source").path("bootstrap-servers").asText();
String topic = root.path("kafka-source").path("topics").get(0).asText();
String groupId = root.path("kafka-source").path("group-id").asText();
// Hand-rolled parsing and fragile validation for business fields
double threshold = root.path("alerting").path("threshold-amount").asDouble();
String email = root.path("alerting").path("notification-email").asText();
```
---
## 2. Declarative YAML schema
Your configuration files map 1:1 to your type-safe DTOs. Instead of scattering parameters across ad-hoc CLI arguments, Java System Properties, and flat properties files, Flinkboot organizes everything into a hierarchical YAML contract supporting profile activation, parameter placeholders, and environment variable substitution.
```yaml
# application.yaml - Maps 1:1 to your AppConfig DTO
job:
name: "order-fraud-detector"
parallelism: 8
checkpointing:
interval: 60000
timeout: 120000
state-backend:
type: "rocksdb"
incremental: true
kafka-source:
name: "fraud-orders-source"
bootstrap-servers:
- "kafka-1.internal.net:9092"
- "kafka-2.internal.net:9092"
topics:
- "orders-v1"
group-id: "fraud-detector-service"
starting-offsets: LATEST
# Universal escape hatch for vendor client tuning & SSL credentials
properties:
security.protocol: "SSL"
ssl.truststore.location: "/var/private/ssl/kafka.truststore.jks"
# Auto-templated at runtime from host or container environment variables
ssl.truststore.password: "${KAFKA_TRUSTSTORE_PASSWORD}"
fetch.max.wait.ms: "500"
alerting:
threshold-amount: 5000.00
notification-email: "fraud-alerts@company.com"
```
```text
# Hardcoded, flat properties or repetitive CLI flags
--job.name order-fraud-detector \
--parallelism 8 \
--checkpoint.interval 60000 \
--rocksdb.incremental true \
--kafka.bootstrap.servers kafka-1.internal.net:9092,kafka-2.internal.net:9092 \
--kafka.topics orders-v1 \
--kafka.group.id fraud-detector-service \
--kafka.properties.security.protocol SSL \
--kafka.properties.ssl.truststore.location /var/private/ssl/kafka.truststore.jks \
--kafka.properties.ssl.truststore.password ${KAFKA_TRUSTSTORE_PASSWORD} \
--kafka.properties.fetch.max.wait.ms 500 \
--alerting.threshold-amount 5000.00 \
--alerting.notification-email fraud-alerts@company.com
# Fragile CLI parsing, lacks hierarchy, and offers zero validation
# when a parameter is misspelled or missing at runtime.
```
---
## 3. Fail-fast application bootstrap
Instead of handwriting hundreds of lines of imperative setup code, Flinkboot bootstraps your entire application in a single statement. Even if you've never configured RocksDB state backends, unaligned checkpoints, latency metrics, or exponential backoff restart strategies before, they are already built-in, pre-tuned for production, and ready to be declared in your YAML with zero plumbing code required.
```java
public static void main(String[] args) throws Exception {
// 1. Initialize Flinkboot with CLI arguments
Flinkboot boot = Flinkboot.initialize(args);
// 2. Load and validate YAML configurations fail-fast into your type-safe record
AppConfig config = boot.configuration(AppConfig.class);
// 3. Pre-configured environment: RocksDB, checkpoints, restarts from YAML in one call
StreamExecutionEnvironment env = boot.executionEnvironment(config.job());
// 4. Run pipeline
env.execute(config.job().name());
}
```
```java
public static void main(String[] args) throws Exception {
ParameterTool params = ParameterTool.fromArgs(args);
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Imperative wiring scattered across the main method
env.setParallelism(params.getInt("parallelism", 1));
env.enableCheckpointing(params.getLong("checkpoint.interval", 60000L));
env.getCheckpointConfig().setCheckpointTimeout(params.getLong("checkpoint.timeout", 120000L));
env.setStateBackend(new EmbeddedRocksDBStateBackend(true));
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.seconds(10)));
// Misconfigurations and type mismatches crash only after job submission
env.execute("LegacyJob");
}
```
---
## 4. Turnkey production connectors
Instantiating sources and sinks in vanilla Flink requires verbose builders, duplicate properties, and manual schema binding. Flinkboot connector factories create fully tuned sources and sinks directly from your validated configuration objects with native serializer resolution.
```java
// Instantiates a fully configured, production-ready KafkaSource in one line
KafkaSource source = KafkaSourceFactory.supplyFor(
config.kafkaSource(),
deserializationSchema
);
```
```java
// Repetitive builder with manual string mappings and custom deserializers
KafkaSource source = KafkaSource.builder()
.setBootstrapServers(brokers)
.setTopics(topic)
.setGroupId(groupId)
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new OrderEventDeserializationSchema())
.setProperty("enable.auto.commit", "false")
.build();
```
---
## 5. Native collection & JDK serialization
In vanilla Flink, collections like `List` inside a POJO silently fall back to Kryo serialization because Flink's type extractor cannot resolve generic parameters. Flinkboot provides turnkey `TypeInfoFactory` classes to guarantee high-throughput, native Flink serializers.
```java
public class OrderEvent {
public String id;
// Instructs Flink to resolve List natively as Types.LIST(Types.STRING)
@TypeInfo(ListTypeInfoFactory.class)
public List tags;
}
```
```java
public class OrderEvent {
public String id;
public List tags; // Flink cannot extract generic parameter E!
}
// In standard Flink, TypeInformation.of(OrderEvent.class) silently assigns
// GenericTypeInfo (Kryo) to the tags field, degrading streaming throughput
// by 3x to 10x and breaking savepoint state schema evolution.
```
---
## 6. Build-time POJO compliance auditing
Apache Flink relies on its high-performance `PojoSerializer` to achieve maximum streaming throughput. If an event class lacks a default constructor, contains an unmapped collection, or misses getters/setters, Flink silently falls back to slow Kryo serialization without failing. Flinkboot provides build-time assertions to guarantee POJO compliance in your unit tests.
```java
@Test
void verifyOrderEventSerialization() {
// Build-time guarantee: recursively audits all fields, getters, and constructors
// Fails the test immediately if any field falls back to Kryo
FlinkbootAssertions.assertThat(OrderEvent.class)
.isPojo();
}
```
```java
// No build-time guarantee in standard Flink!
// Developers either discover severe performance degradation in production,
// or attempt brittle runtime TypeInformation inspection:
TypeInformation ti = TypeInformation.of(OrderEvent.class);
// Returns GenericTypeInfo silently in production when POJO rules are violated,
// slowing down streaming pipelines by 3x to 10x without any explicit error.
```
================================================================================
DOCUMENT: Bootstrapping execution environment (/docs/configuration/auto-configure-execution-environment)
================================================================================
# 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.
```java
// 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:
```yaml
# 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 Key | Type | Required | Description |
|:--------------|:-------|:---------|:------------|
| `name` | String | **Yes** | Canonical job name registered with Flink (`PipelineOptions.NAME`). |
| `environment` | Object | No | Execution environment settings (`ExecutionEnvironmentProperties`). |
---
### B. Execution options (`environment.execution:`)
| Property Key | Type | Required | Validation | Description |
|:--------------------------|:-----------|:---------|:----------------------------|:------------|
| `runtime-mode` | Enum | No | Enum | Execution runtime mode: `STREAMING`, `BATCH`, or `AUTOMATIC` (`ExecutionOptions.RUNTIME_MODE`). |
| `parallelism` | Integer | No | `@Positive` | Default execution parallelism (`CoreOptions.DEFAULT_PARALLELISM`). |
| `max-parallelism` | Integer | No | `@Positive` | Maximum parallelism for key groups rescale (`PipelineOptions.MAX_PARALLELISM`). |
| `buffer-timeout` | `Duration` | No | `@DurationMin(millis = 0)` | Buffer timeout (`ExecutionOptions.BUFFER_TIMEOUT`), e.g. `"PT0.05S"`. |
| `auto-watermark-interval` | `Duration` | No | `@DurationMin(millis = 0)` | Periodic watermark emission interval (`PipelineOptions.AUTO_WATERMARK_INTERVAL`), e.g. `"PT0.2S"`. |
| `object-reuse` | Boolean | No | Boolean | Enable object reuse optimization (`PipelineOptions.OBJECT_REUSE`). Defaults to `false`. |
---
### C. Fault tolerance & checkpointing (`environment.checkpointing:`)
| Property Key | Type | Required | Validation | Description |
|:----------------------------------|:-----------|:---------|:----------------------------|:------------|
| `enabled` | Boolean | No | Boolean | Master switch for checkpointing. Defaults to `true` when block is present. |
| `interval` | `Duration` | No | `@DurationMin(millis = 1)` | Time interval between checkpoints (`CheckpointingOptions.CHECKPOINTING_INTERVAL`), e.g. `"PT30S"`. |
| `mode` | Enum | No | Enum | Consistency mode: `EXACTLY_ONCE` or `AT_LEAST_ONCE`. |
| `timeout` | `Duration` | No | `@DurationMin(millis = 1)` | Maximum checkpoint duration before aborting (`CheckpointingOptions.CHECKPOINTING_TIMEOUT`). |
| `min-pause-between-checkpoints` | `Duration` | No | `@DurationMin(millis = 0)` | Minimum rest duration between consecutive checkpoints. |
| `max-concurrent-checkpoints` | Integer | No | `@Positive` | Maximum concurrent checkpoints allowed. |
| `externalized-checkpoint-cleanup` | Enum | No | Enum | Cleanup retention on cancel: `RETAIN_ON_CANCELLATION`, `DELETE_ON_CANCELLATION`, or `NO_EXTERNALIZED_CHECKPOINTS`. |
| `unaligned-checkpoints` | Boolean | No | Boolean | Enable unaligned checkpoints (`CheckpointingOptions.ENABLE_UNALIGNED`). |
| `aligned-checkpoint-timeout` | `Duration` | No | `@DurationMin(millis = 0)` | Timeout before switching to unaligned checkpoints, e.g. `"PT1S"`. |
| `storage-uri` | String | No | String | Target checkpoint storage URI, e.g. `s3://bucket/checkpoints`. |
---
### D. State backend & RocksDB (`environment.state-backend:`)
| Property Key | Type | Required | Validation | Description |
|:---------------------|:--------|:------------------------------|:-----------|:------------|
| `type` | Enum | No | Enum | State backend type: `ROCKSDB`, `HASHMAP`, `CHANGELOG`, or `CUSTOM`. |
| `checkpoint-storage` | Enum | No | Enum | Storage mechanism: `JOBMANAGER` or `FILESYSTEM`. |
| `incremental` | Boolean | No | Boolean | Enable incremental checkpoints for RocksDB (`CheckpointingOptions.INCREMENTAL_CHECKPOINTS`). |
| `latency-tracking` | Boolean | No | Boolean | Enable latency tracking metrics for state access (`StateBackendOptions.LATENCY_TRACK_ENABLED`). |
| `custom-class` | String | **Yes** (if `type == CUSTOM`) | String | Fully 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 Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `type` | Enum | No | Enum | Strategy type: `NO_RESTART`, `FIXED_DELAY`, `FAILURE_RATE`, `EXPONENTIAL_DELAY`, or `FALLBACK`. |
| `fixed-delay` | Object | No | `@Valid` | Parameters for `FIXED_DELAY` strategy. |
| `failure-rate` | Object | No | `@Valid` | Parameters for `FAILURE_RATE` strategy. |
| `exponential-delay` | Object | No | `@Valid` | Parameters for `EXPONENTIAL_DELAY` strategy. |
#### 1. Fixed delay (`type: FIXED_DELAY`)
```yaml
restart-strategy:
type: "FIXED_DELAY"
fixed-delay:
attempts: 3
delay: "PT10S"
```
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `attempts` | Integer | **Yes** | `@PositiveOrZero` | Number of restart attempts before job failure. |
| `delay` | `Duration` | No | `@DurationMin(millis = 0)` | Delay between restart attempts. Defaults to `0s`. |
#### 2. Failure rate (`type: FAILURE_RATE`)
```yaml
restart-strategy:
type: "FAILURE_RATE"
failure-rate:
max-failures-per-interval: 5
failure-interval: "PT5M"
delay: "PT10S"
```
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `max-failures-per-interval` | Integer | **Yes** | `@Positive` | Maximum failures permitted within the time window. |
| `failure-interval` | `Duration` | **Yes** | `@DurationMin(millis = 1)` | Measurement window interval. |
| `delay` | `Duration` | No | `@DurationMin(millis = 0)` | Delay between restart attempts. Defaults to `0s`. |
#### 3. Exponential delay (`type: EXPONENTIAL_DELAY`)
```yaml
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 Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `initial-backoff` | `Duration` | **Yes** | `@DurationMin(millis = 1)` | Initial backoff duration. |
| `max-backoff` | `Duration` | **Yes** | Greater than or equal to `initial-backoff` | Maximum cap for backoff duration. |
| `backoff-multiplier` | Double | No | Strictly `> 1.0` | Exponential multiplier for consecutive failures. Defaults to `2.0`. |
| `reset-backoff-threshold` | `Duration` | No | `@DurationMin(millis = 1)` | Failure-free uptime required to reset backoff. Defaults to `PT1H`. |
| `jitter-factor` | Double | No | `0.0 <= x <= 1.0` | Random jitter ratio added to backoff. Defaults to `0.1`. |
---
### F. Savepoint recovery (`environment.savepoint-restore:`)
```yaml
savepoint-restore:
savepoint-path: "/mnt/savepoints/savepoint-0001"
allow-non-restored-state: false
restore-mode: "CLAIM"
```
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `savepoint-path` | String | **Yes** | `@NotBlank` | Path to savepoint or initial checkpoint directory. |
| `allow-non-restored-state` | Boolean | No | Boolean | Start even if savepoint contains unmapped subtask state. |
| `restore-mode` | Enum | No | Enum | Restore mode: `CLAIM`, `NO_CLAIM`, or `LEGACY`. |
---
### G. Local dev WebUI (`environment.local-web-ui:`)
```yaml
local-web-ui:
enabled: true
port: 8081
bind-address: "localhost"
```
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `enabled` | Boolean | **Yes** (if block present) | `@NotNull` | Starts a local Flink MiniCluster with the WebUI dashboard active during IDE testing. |
| `port` | Integer | No | `@Range(min = 0, max = 65535)` | WebUI REST port. Defaults to `8081`. Set to `0` for dynamic port allocation. |
| `bind-address` | String | No | `@NotBlank` | WebUI host bind address. Defaults to `localhost`. |
Enabling `local-web-ui.enabled: true` requires `org.apache.flink:flink-runtime-web` on the classpath:
```xml
org.apache.flink
flink-runtime-web
provided
```
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:
```yaml
environment:
properties:
taskmanager.memory.managed.fraction: "0.4"
pipeline.operator-chaining.enabled: "true"
```
| Property Key | Type | Required | Description |
|:---|:---|:---|:---|
| `properties` | `Map` | No | Free-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. |
================================================================================
DOCUMENT: Loading configuration (/docs/configuration/loading-configuration)
================================================================================
# Loading configuration
In real-world Apache Flink applications, pipelines require custom business parameters alongside execution settings: alert thresholds, window intervals, database endpoints, or external API keys.
Under the hood, configuration deserialization is powered by [Jackson](https://github.com/FasterXML/jackson). Standard Java types, collections, and Java temporal types are supported out of the box through Jackson's Java Time module.
YAML keys map to Java fields using Jackson's `@JsonProperty("key-name")` annotation. Property names are matched case-insensitively.
---
## 1. Defining configuration models
### Mandatory fields with Java Records (Java 17+)
For models where all fields are required and you are running on Java 17+, standard Java Records provide clean, compact immutability:
```java
package com.company.fraud.config;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Positive;
import java.io.Serializable;
import java.time.Duration;
public record AlertingProperties(
@NotNull @Positive @JsonProperty("threshold-amount") Double thresholdAmount,
@NotNull @JsonProperty("evaluation-window") Duration evaluationWindow,
@NotBlank @JsonProperty("notification-email") String notificationEmail
) implements Serializable {}
```
Corresponding YAML snippet:
```yaml
alerting:
threshold-amount: 5000.00
evaluation-window: "PT5M"
notification-email: "fraud-ops@company.com"
```
### Immutable class with Optional getters for optional fields
When certain configuration properties are optional, a Java Record cannot safely hide the raw nullable accessor. Instead, define an immutable class with private final fields, `@JsonCreator`, and explicit `Optional` getters:
```java
package com.company.fraud.config;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.validation.constraints.NotBlank;
import java.io.Serializable;
import java.util.Optional;
public final class ServerProperties implements Serializable {
private final String host;
private final Integer port;
@JsonCreator
public ServerProperties(
@NotBlank @JsonProperty("host") String host,
@JsonProperty("port") Integer port
) {
this.host = host;
this.port = port;
}
public String host() {
return host;
}
public Optional port() {
return Optional.ofNullable(port);
}
}
```
Corresponding YAML snippet:
```yaml
server:
host: "db.internal.net"
# port is omitted and resolves to Optional.empty()
```
---
## 2. Customizing Jackson YAML deserialization
If your domain models require custom Jackson configuration or third-party modules, supply a builder customizer lambda when loading the configuration:
```java
AppConfig config = boot.configuration(AppConfig.class, builder -> {
builder.configure(DeserializationFeature.READ_UNKNOWN_ENUM_VALUES_AS_NULL, true);
builder.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, true);
});
```
You can also pass a pre-configured `YAMLMapper` directly:
```java
YAMLMapper mapper = new YAMLMapper();
// Custom mapper configuration...
AppConfig config = boot.configuration(AppConfig.class, mapper);
```
For more options, refer to the official [Jackson documentation](https://github.com/FasterXML/jackson-databind).
---
## 3. Loading configuration
To bind your YAML files to your domain model, call `boot.configuration(Class)`:
```java
public static void main(String[] args) throws Exception {
Flinkboot boot = Flinkboot.initialize(args);
// Binds and validates YAML configuration into your model
AppConfig config = boot.configuration(AppConfig.class);
// Access your strongly-typed configuration
double threshold = config.alerting().thresholdAmount();
String host = config.server().host();
int port = config.server().port().orElse(8080);
}
```
### Default configuration location
By default, Flinkboot looks for a file named `job-configuration.yaml` in your application classpath:
```text
classpath:job-configuration.yaml
```
If this file is missing and no other location is specified, initialization fails fast.
### Overriding configuration paths
You can override the default location at runtime or supply multiple configuration files to merge:
* **Via Command Line Argument:** Pass `-flinkboot-configurations` followed by comma-separated paths:
```bash
flink run MyJob.jar -flinkboot-configurations "file:/etc/flink/job-config.yaml"
```
To merge multiple configurations sequentially:
```bash
flink run MyJob.jar -flinkboot-configurations "classpath:base-config.yaml,file:/etc/flink/production.yaml"
```
* **Via Environment Variable:**
```bash
export FLINKBOOT_CONFIGURATIONS="file:/etc/flink/production.yaml"
flink run MyJob.jar
```
Supported URI schemes include `classpath:`, `resource:`, and `file:`.
### Merging multiple configuration files
When loading multiple files sequentially, Flinkboot enforces strict collision rules by default to prevent accidental configuration overwrites:
* **Scalar Overrides (`--flinkboot-configuration-override`):**
By default (`false`), redefining an already existing key across files halts startup immediately with a `YamlParsingException`. To intentionally allow later files to overwrite earlier scalar values (e.g. environment-specific overrides over a base template), enable:
```bash
flink run MyJob.jar \
-flinkboot-configurations "classpath:base.yaml,file:/etc/flink/prod.yaml" \
--flinkboot-configuration-override
```
*(Or set environment variable `FLINKBOOT_CONFIGURATION_OVERRIDE=true`)*.
* **List Merging (`--flinkboot-configuration-list-merging`):**
By default (`false`), redefining an existing list or array across merged files is considered a collision and throws an exception. To append and concatenate list items from subsequent files, enable:
```bash
flink run MyJob.jar \
-flinkboot-configurations "classpath:base.yaml,file:/etc/flink/prod.yaml" \
--flinkboot-configuration-list-merging
```
*(Or set environment variable `FLINKBOOT_CONFIGURATION_LIST_MERGING=true`)*.
---
## 4. Environment variable templating
Configuration files support dynamic variable interpolation using the `${VAR_NAME}` syntax. Values are resolved from host environment variables before Jackson binds your configuration models:
```yaml
server:
host: "${DATABASE_HOST}"
port: ${DATABASE_PORT}
```
### Fail-fast invariant
Flinkboot validates all placeholders strictly at startup:
- If a referenced variable is missing from the host environment, initialization fails fast with an `UnresolvedPropertyPlaceholderException`.
- This ensures your streaming job never starts with partially resolved or missing credentials.
### Escaping literal placeholders
If your configuration contains literal `${...}` strings that should not be evaluated as environment variables (for example, partitioned file patterns or regex templates), escape the prefix with a backslash:
```yaml
sink:
path-template: "\${year}/\${month}/\${day}"
```
Flinkboot preserves the literal pattern `${year}/${month}/${day}` without attempting resolution.
================================================================================
DOCUMENT: Resolving CLI & environment parameters (/docs/configuration/resolving-cli-and-environment-parameters)
================================================================================
# Resolving CLI & environment parameters
In addition to static YAML files, streaming applications frequently need to inspect runtime command-line arguments, toggle operational flags, read external resources, and leverage framework-level flags.
Flinkboot unifies CLI arguments and environment variables into a single, fail-fast evaluation engine with automatic case normalization and strict collision safeguards.
---
## 1. Custom parameters (`boot.parameter`)
A **parameter** is a key-value pair where the value is a string. Parameters are retrieved as a Java `Optional`:
```java
Flinkboot boot = Flinkboot.initialize(args);
// Retrieves optional parameter "db-url"
Optional dbUrl = boot.parameter("db-url");
String connectionString = dbUrl.orElse("jdbc:postgresql://localhost:5432/defaultdb");
```
### Passing parameters
* **Via Command Line (CLI):** Prefix the parameter key with a single dash `-`:
```bash
flink run MyJob.jar -db-url "jdbc:postgresql://prod-db:5432/orders"
```
* **Via Environment Variable:** Flinkboot normalizes keys to uppercase and replaces dashes with underscores:
```bash
export DB_URL="jdbc:postgresql://prod-db:5432/orders"
flink run MyJob.jar
```
### Resolution precedence
1. **CLI Argument (`-db-url`)**: Highest priority.
2. **Environment Variable (`DB_URL`)**: Evaluated if CLI argument is absent.
3. **Empty Optional**: Returns `Optional.empty()` if defined in neither.
---
## 2. Boolean flags (`boot.flag`)
A **flag** is a boolean switch used to enable or disable runtime behaviors (debug modes, dry runs, backfills). By default, flags evaluate to `false` when omitted:
```java
Flinkboot boot = Flinkboot.initialize(args);
// Evaluates boolean flag "dry-run"
boolean isDryRun = boot.flag("dry-run");
if (isDryRun) {
System.out.println("Executing in dry-run mode without committing state.");
}
```
### Passing flags
* **Via Command Line (CLI):** Prefix the flag key with double dashes `--`:
```bash
flink run MyJob.jar --dry-run
```
*(The presence of `--dry-run` automatically sets the flag to `true`)*.
* **Via Environment Variable:**
```bash
export DRY_RUN=true
flink run MyJob.jar
```
### Strict boolean parsing
When passing flags via environment variables, values must strictly be `"true"` or `"false"` (case-insensitive). Any arbitrary string (such as `"yes"`, `"1"`, or `"on"`) fails fast at startup with a `BooleanParsingException`.
---
## 3. Built-in values checked by Flinkboot
Flinkboot provides built-in options that control configuration resolution, merging, and validation. These keys are reserved by the framework:
| Command Line Key | Environment Variable | Default | Purpose & Fail-Fast Check |
|:---|:---|:---|:---|
| `-flinkboot-configurations ` | `FLINKBOOT_CONFIGURATIONS` | `classpath:job-configuration.yaml` | Comma-separated list of configuration URIs to load and merge sequentially. |
| `--flinkboot-configuration-override` | `FLINKBOOT_CONFIGURATION_OVERRIDE` | `false` | When `false`, redefining an existing scalar key across merged files throws a `YamlParsingException`. When `true`, overrides are allowed. |
| `--flinkboot-configuration-list-merging` | `FLINKBOOT_CONFIGURATION_LIST_MERGING` | `false` | When `false`, redefining a list throws an exception. When `true`, list items across files are concatenated. |
| `--flinkboot-configuration-disable-validation` | `FLINKBOOT_CONFIGURATION_DISABLE_VALIDATION` | `false` | Bypasses Jakarta Bean Validation checks on configuration models (useful for mock tests). Malformed YAML still fails fast. |
| `-flinkboot-configuration-violations-log-size ` | `FLINKBOOT_CONFIGURATION_VIOLATIONS_LOG_SIZE` | `10` | Maximum number of Bean Validation violations displayed before summary truncation. Must be a **strictly positive integer**; zero or negative values fail fast at startup. |
### Automatic framework safeguards
**Reserved key collision protection:**
You cannot name a custom parameter `-flinkboot-configurations` or custom flag `--flinkboot-configuration-override`. Flinkboot prevents accidental interception of internal framework flags.
================================================================================
DOCUMENT: Validating configuration (/docs/configuration/validating-configuration)
================================================================================
# Validating configuration
Configuration errors are a frequent cause of streaming runtime failures, often surfacing only after job submission. Flinkboot integrates [Jakarta Bean Validation](https://jakarta.ee/specifications/bean-validation/) to validate models fail-fast during initialization before your streaming pipeline starts.
---
## 1. Jakarta Bean Validation
Validation occurs **after the final configuration tree has been fully resolved and merged** (including all configuration files and environment variable substitutions). Every model loaded via `boot.configuration(...)` is then automatically evaluated against standard Jakarta Bean Validation constraints.
You can apply any constraint from the `jakarta.validation.constraints` package, such as `@NotNull`, `@NotBlank`, `@Positive`, `@Min`, `@Max`, `@Pattern`, or `@Size`.
```java
package com.company.fraud.config;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import jakarta.validation.constraints.Positive;
import java.io.Serializable;
public record DatabaseProperties(
@NotBlank @JsonProperty("host") String host,
@Min(1024) @Max(65535) @JsonProperty("port") int port,
@NotNull @Positive @JsonProperty("max-connections") Integer maxConnections
) implements Serializable {}
```
For nested models, remember to place `@Valid` on the parent field so that Jakarta cascades validation to inner objects:
```java
public record AppConfig(
@Valid @NotNull @JsonProperty("database") DatabaseProperties database
) implements Serializable {}
```
Refer to the official [Jakarta Bean Validation documentation](https://jakarta.ee/specifications/bean-validation/) for details on all available built-in annotations.
---
## 2. Validation error reporting
When a constraint is violated, Flinkboot halts application startup with a `ConfigurationValidationException`. The violation path lists the model property paths so the developer can immediately locate the offending configuration field.
Example log output:
```text
io.github.sekelenao.flinkboot.core.api.exception.configuration.ConfigurationValidationException:
Configuration validation failed with 2 violation(s):
- database.port: must be greater than or equal to 1024
- database.maxConnections: must be greater than 0
```
By default, Flinkboot prints up to 10 validation violations before summarizing any remaining errors (`... and X more violation(s)`). You can change this limit using the CLI parameter:
```bash
-flinkboot-configuration-violations-log-size 25
```
---
## 3. Disabling validation
Bean validation can be disabled using the CLI flag `--flinkboot-configuration-disable-validation`.
However, disabling validation means no constraints are verified at startup. No application behavior or runtime stability is guaranteed when running with validation disabled.
---
## 4. Cross-field validation with `ValidatableProperties`
Field-level annotations cannot validate invariants that span multiple fields, such as ensuring a sliding window slide interval is strictly shorter than its window length.
Flinkboot provides the `ValidatableProperties` interface to implement custom cross-field validation rules directly on your model:
```java
package com.company.fraud.config;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.github.sekelenao.flinkboot.core.api.validation.ValidatableProperties;
import jakarta.validation.ConstraintValidatorContext;
import jakarta.validation.constraints.NotNull;
import java.io.Serializable;
import java.time.Duration;
public record WindowProperties(
@NotNull @JsonProperty("window-size") Duration windowSize,
@JsonProperty("slide-duration") Duration slideDuration
) implements ValidatableProperties, Serializable {
@Override
public boolean validate(ConstraintValidatorContext context) {
if (slideDuration != null && slideDuration.compareTo(windowSize) >= 0) {
// Suppress the default generic class-level message
context.disableDefaultConstraintViolation();
// Bind the violation to the specific field in error
context.buildConstraintViolationWithTemplate(
"slide-duration must be strictly less than window-size")
.addPropertyNode("slide-duration")
.addConstraintViolation();
return false;
}
return true;
}
}
```
Calling `context.disableDefaultConstraintViolation()` suppresses generic top-level class messages and outputs your specific violation message targeting the relevant property node.
================================================================================
DOCUMENT: Apache Fluss (/docs/connectors/fluss)
================================================================================
# Apache Fluss
Flinkboot provides typed configuration models and factories to initialize Apache Flink's `FlussSource` and `FlussSink` directly from declarative YAML configurations.
---
## 1. Maven dependency
Add `flinkboot-fluss` to your `pom.xml`. Versions are managed automatically by the Flinkboot BOM:
```xml
io.github.sekelenao
flinkboot-fluss
```
---
## 2. Fluss source
### YAML configuration
```yaml
fluss-source:
name: "events-source"
bootstrap-servers:
- "localhost:9123"
database: "analytics_db"
table: "user_events"
startup-mode: "EARLIEST"
properties:
client.scanner.fetch.max-bytes: "1048576"
```
For timestamp-based positioning:
```yaml
fluss-source:
name: "replay-source"
bootstrap-servers:
- "localhost:9123"
database: "analytics_db"
table: "user_events"
startup-mode: "TIMESTAMP"
startup-timestamp: 1700000000000
```
### Configuration reference
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `name` | String | **Yes** | `@NotBlank` | Unique operator identifier in the Flink DAG graph. |
| `bootstrap-servers` | `List` | **Yes** | `@NotEmpty`, items `@NotBlank` | Addresses of Fluss coordinators. |
| `database` | String | **Yes** | `@NotBlank` | Target Fluss database. |
| `table` | String | **Yes** | `@NotBlank` | Target Fluss table. |
| `startup-mode` | Enum | **Yes** | `@NotNull` | Startup strategy: `EARLIEST`, `LATEST`, `FULL`, `TIMESTAMP`. |
| `startup-timestamp` | Long | Conditional | `@PositiveOrZero` | Epoch millisecond timestamp (**mandatory** if `startup-mode: TIMESTAMP`, forbidden otherwise). |
| `properties` | `Map` | No | `@NotNull` entries | Additional Fluss client/scanner tuning options. |
---
## 3. Fluss sink
### YAML configuration
```yaml
fluss-sink:
name: "aggregates-sink"
bootstrap-servers:
- "localhost:9123"
database: "analytics_db"
table: "user_aggregates"
properties:
client.writer.batch-size: "1mb"
client.writer.batch-timeout: "50ms"
```
### Configuration reference
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `name` | String | **Yes** | `@NotBlank` | Unique operator identifier in the Flink DAG graph. |
| `bootstrap-servers` | `List` | **Yes** | `@NotEmpty`, items `@NotBlank` | Addresses of Fluss coordinators. |
| `database` | String | **Yes** | `@NotBlank` | Target Fluss database. |
| `table` | String | **Yes** | `@NotBlank` | Target Fluss table. |
| `properties` | `Map` | No | `@NotNull` entries | Additional Fluss client/writer tuning options. |
---
## 4. Pipeline integration
Combine `FlussSourceProperties` and `FlussSinkProperties` in your configuration model and initialize them with `FlussSourceFactory` and `FlussSinkFactory`:
```java
package com.company.analytics.config;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.github.sekelenao.flinkboot.core.api.properties.JobProperties;
import io.github.sekelenao.flinkboot.fluss.api.properties.sink.FlussSinkProperties;
import io.github.sekelenao.flinkboot.fluss.api.properties.source.FlussSourceProperties;
import jakarta.validation.Valid;
import jakarta.validation.constraints.NotNull;
import java.io.Serializable;
public record AppConfig(
@Valid @NotNull @JsonProperty("job") JobProperties job,
@Valid @NotNull @JsonProperty("fluss-source") FlussSourceProperties flussSource,
@Valid @NotNull @JsonProperty("fluss-sink") FlussSinkProperties flussSink
) implements Serializable {}
```
In your main entrypoint:
```java
package com.company.analytics;
import com.company.analytics.config.AppConfig;
import io.github.sekelenao.flinkboot.core.api.Flinkboot;
import io.github.sekelenao.flinkboot.fluss.api.sink.FlussSinkFactory;
import io.github.sekelenao.flinkboot.fluss.api.source.FlussSourceFactory;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.data.RowData;
import org.apache.fluss.flink.sink.FlussSink;
import org.apache.fluss.flink.sink.serializer.RowDataSerializationSchema;
import org.apache.fluss.flink.source.FlussSource;
import org.apache.fluss.flink.source.deserializer.RowDataDeserializationSchema;
public class FlussPipelineJob {
public static void main(String[] args) throws Exception {
Flinkboot boot = Flinkboot.initialize(args);
AppConfig config = boot.configuration(AppConfig.class);
StreamExecutionEnvironment env = boot.executionEnvironment(config.job());
// 1. Build Fluss Source
RowDataDeserializationSchema deserializer = ...; // Your deserialization schema
FlussSource source = FlussSourceFactory.supplyFor(config.flussSource(), deserializer);
// 2. Build Fluss Sink
RowDataSerializationSchema serializer = ...; // Your serialization schema
FlussSink sink = FlussSinkFactory.supplyFor(config.flussSink(), serializer);
// 3. Connect stream pipeline
DataStream stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
config.flussSource().name()
);
stream.sinkTo(sink).name(config.flussSink().name());
env.execute(config.job().name());
}
}
```
---
## 5. Java 17+ and Apache Arrow JVM options
Apache Fluss utilizes Apache Arrow for off-heap buffer management. On Java 17 and later, add the following JVM argument to open `java.nio` to unnamed modules:
```bash
--add-opens=java.base/java.nio=ALL-UNNAMED
```
When running unit/integration tests with Maven Surefire:
```xml
org.apache.maven.plugins
maven-surefire-plugin
--add-opens=java.base/java.nio=ALL-UNNAMED
```
================================================================================
DOCUMENT: Apache Kafka (/docs/connectors/kafka)
================================================================================
# Apache Kafka
Flinkboot provides typed configuration models and factories to initialize Apache Flink's `KafkaSource` and `KafkaSink` directly from declarative YAML configurations.
---
## 1. Maven dependency
Add `flinkboot-kafka` to your `pom.xml`. Versions are managed automatically by the Flinkboot BOM:
```xml
io.github.sekelenao
flinkboot-kafka
```
---
## 2. Kafka source
### YAML configuration
You can configure subscriptions using either an explicit topic list or a regex pattern:
```yaml
kafka-source:
name: "orders-source"
bootstrap-servers:
- "localhost:9092"
group-id: "order-consumers"
topics:
- "orders"
- "payments"
starting-offsets: "EARLIEST"
properties:
session.timeout.ms: "45000"
```
For timestamp-based positioning:
```yaml
kafka-source:
name: "replay-orders-source"
bootstrap-servers:
- "localhost:9092"
group-id: "replay-consumers"
topics:
- "orders"
starting-offsets: "TIMESTAMP"
starting-offsets-timestamp: 1689717600000
```
For explicit partition offsets:
```yaml
kafka-source:
name: "partition-orders-source"
bootstrap-servers:
- "localhost:9092"
group-id: "partition-consumers"
topics:
- "orders"
starting-offsets: "OFFSETS"
starting-offsets-partition-offsets:
- topic: "orders"
partition: 0
offset: 12500
- topic: "orders"
partition: 1
offset: 14200
```
### Configuration reference
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `name` | String | **Yes** | `@NotBlank` | Operator name registered in the Flink DAG graph. |
| `bootstrap-servers` | `List` | **Yes** | `@NotEmpty`, items `@NotBlank` | Kafka bootstrap broker hosts and ports. |
| `group-id` | String | **Yes** | `@NotBlank` | Consumer group ID. |
| `topics` | `List` | Conditional | items `@NotBlank` | Explicit topic subscriptions (mutually exclusive with `topic-pattern`). |
| `topic-pattern` | String | Conditional | Valid regex | Topic subscription regex pattern (mutually exclusive with `topics`). |
| `starting-offsets` | Enum | **Yes** | `@NotNull` | Startup strategy: `EARLIEST`, `LATEST`, `COMMITTED`, `TIMESTAMP`, `OFFSETS`. |
| `starting-offsets-timestamp` | Long | Conditional | `@PositiveOrZero` | Timestamp in milliseconds (**mandatory** if `starting-offsets: TIMESTAMP`, forbidden otherwise). |
| `starting-offsets-partition-offsets` | List | Conditional | `@Valid` items | List of partition starting offsets (**mandatory** if `starting-offsets: OFFSETS`, forbidden otherwise). |
| `properties` | `Map` | No | Free-form map | Additional Kafka consumer tuning properties (e.g. `session.timeout.ms`). |
### Offset strategies (`starting-offsets`)
| Strategy | Required Complementary Keys | Forbidden Complementary Keys | Description |
|:---|:---|:---|:---|
| `EARLIEST` | None | `starting-offsets-timestamp`, `starting-offsets-partition-offsets` | Start from earliest available log offsets. |
| `LATEST` | None | `starting-offsets-timestamp`, `starting-offsets-partition-offsets` | Start from latest log offsets. |
| `COMMITTED` | None | `starting-offsets-timestamp`, `starting-offsets-partition-offsets` | Start from consumer group committed offsets. |
| `TIMESTAMP` | `starting-offsets-timestamp` (`Long`) | `starting-offsets-partition-offsets` | Position based on record epoch millisecond timestamps. |
| `OFFSETS` | `starting-offsets-partition-offsets` (`List`) | `starting-offsets-timestamp` | Explicit starting offsets per topic partition. |
---
## 3. Kafka sink
### YAML configuration
```yaml
kafka-sink:
name: "alerts-sink"
bootstrap-servers:
- "localhost:9092"
topic: "fraud-alerts"
delivery-guarantee: "EXACTLY_ONCE"
transactional-id-prefix: "fraud-evaluator"
properties:
acks: "all"
```
### Configuration reference
| Property Key | Type | Required | Validation | Description |
|:---|:---|:---|:---|:---|
| `name` | String | **Yes** | `@NotBlank` | Logical identifier for the sink operator. |
| `bootstrap-servers` | `List` | **Yes** | `@NotEmpty`, items `@NotBlank` | Kafka broker endpoints. |
| `topic` | String | **Yes** | `@NotBlank` | Target Kafka topic for emitted events. |
| `delivery-guarantee` | Enum | **Yes** | `NONE`, `AT_LEAST_ONCE`, `EXACTLY_ONCE` | Delivery semantic guarantee. |
| `transactional-id-prefix` | String | Conditional | String | Transactional prefix. **Mandatory** if `delivery-guarantee` is `EXACTLY_ONCE`, prohibited otherwise. |
| `properties` | `Map` | No | Non-blank keys/values | Custom Kafka producer client settings. |
---
## 4. Pipeline integration
Combine `KafkaSourceProperties` and `KafkaSinkProperties` in your configuration model and instantiate them via `KafkaSourceFactory` and `KafkaSinkFactory`:
```java
package com.company.fraud.config;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.github.sekelenao.flinkboot.core.api.properties.JobProperties;
import io.github.sekelenao.flinkboot.kafka.api.properties.sink.KafkaSinkProperties;
import io.github.sekelenao.flinkboot.kafka.api.properties.source.KafkaSourceProperties;
import jakarta.validation.Valid;
import jakarta.validation.constraints.NotNull;
import java.io.Serializable;
public record AppConfig(
@Valid @NotNull @JsonProperty("job") JobProperties job,
@Valid @NotNull @JsonProperty("kafka-source") KafkaSourceProperties kafkaSource,
@Valid @NotNull @JsonProperty("kafka-sink") KafkaSinkProperties kafkaSink
) implements Serializable {}
```
In your main entrypoint:
```java
package com.company.fraud;
import com.company.fraud.config.AppConfig;
import io.github.sekelenao.flinkboot.core.api.Flinkboot;
import io.github.sekelenao.flinkboot.kafka.api.sink.KafkaSinkFactory;
import io.github.sekelenao.flinkboot.kafka.api.source.KafkaSourceFactory;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class FraudDetectionJob {
public static void main(String[] args) throws Exception {
Flinkboot boot = Flinkboot.initialize(args);
AppConfig config = boot.configuration(AppConfig.class);
StreamExecutionEnvironment env = boot.executionEnvironment(config.job());
// 1. Build Kafka Source
KafkaSource source = KafkaSourceFactory.supplyFor(
config.kafkaSource(),
KafkaRecordDeserializationSchema.valueOnly(new SimpleStringSchema())
);
// 2. Build Kafka Sink
KafkaSink sink = KafkaSinkFactory.supplyFor(
config.kafkaSink(),
KafkaRecordSerializationSchema.builder()
.setTopic(config.kafkaSink().topic())
.setValueSerializationSchema(new SimpleStringSchema())
.build()
);
// 3. Connect stream pipeline
env.fromSource(source, WatermarkStrategy.noWatermarks(), config.kafkaSource().name())
.sinkTo(sink)
.name(config.kafkaSink().name());
env.execute(config.job().name());
}
}
```
If you need programmatic customization on Flink's native builders, use `KafkaSourceFactory.supplyBuilderFor(...)` or `KafkaSinkFactory.supplyBuilderFor(...)`.
================================================================================
DOCUMENT: Configuration serialization for operators (/docs/good-practices/configuration-serialization-for-operators)
================================================================================
# Configuration serialization for operators
In Apache Flink streaming topologies, user-defined functions and operators (`ProcessFunction`, `MapFunction`, sinks, etc.) are instantiated on the client or JobManager and serialized over the network to remote TaskManagers.
When operators require configuration values (such as thresholds, batch sizes, or connection options), passing non-serializable objects causes job submission to fail immediately with a fatal `java.io.NotSerializableException`.
A recommended practice is to load the application configuration in a unit test and assert its Java serialization compliance before shipping fields to operators.
---
## 1. Why assert configuration serialization?
Configuration classes mapped by Flinkboot often hold structured domain properties used across the streaming topology. Flink uses standard Java serialization (`ObjectOutputStream` / `ObjectInputStream`) to transfer operator instances across cluster nodes.
Testing configuration serialization guarantees that:
* The configuration class and all nested properties implement `java.io.Serializable`.
* Any piece of configuration passed as a constructor argument to an operator can be safely distributed across TaskManagers.
* Potential issues are caught early in the test suite rather than during deployment on the cluster.
---
## 2. Testing configuration serialization
Calling `Flinkboot.initialize()` with no arguments loads the default configuration from `src/test/resources/job-configuration.yaml` in the test classpath.
All configuration loading features remain available in tests: you can specify custom paths via `-flinkboot-configurations`, merge multiple configuration files (e.g. `classpath:base.yaml,classpath:override.yaml`), pass CLI overrides, or use environment placeholders.
Once initialized, verify that the configuration instance passes the serialization round-trip using `FlinkbootAssertions.assertThat(...).isSerializable()`:
```java
import io.github.sekelenao.flinkboot.core.api.Flinkboot;
import io.github.sekelenao.flinkboot.test.api.assertion.FlinkbootAssertions;
import org.junit.jupiter.api.Test;
class JobConfigurationTest {
@Test
void shouldBeSerializable() throws Exception {
MyApplicationConfig config = Flinkboot.initialize()
.configuration(MyApplicationConfig.class);
FlinkbootAssertions.assertThat(config)
.isSerializable();
}
}
```
---
## 3. Passing configuration slices to operators
Instead of passing the entire global configuration object to every operator, extract only the necessary nested configuration records or sub-properties:
```java
import org.apache.flink.api.common.functions.MapFunction;
import java.util.Objects;
public class FilterFunction implements MapFunction {
private final FilterConfig config;
public FilterFunction(FilterConfig config) {
this.config = Objects.requireNonNull(config);
}
@Override
public Event map(Event value) {
return value.score() >= config.threshold() ? value : null;
}
}
```
Because `JobConfigurationTest` validates the complete configuration object graph recursively, you can safely pass `config.filter()` or other configuration sections directly into operator constructors without risking runtime serialization errors.
================================================================================
DOCUMENT: Compatibility (/docs/setup/compatibility)
================================================================================
# Compatibility
Supported Apache Flink and Java versions for Flinkboot.
| Flinkboot Version | Apache Flink Version | Java Runtime | Bytecode Target |
|:------------------|:---------------------|:-------------|:----------------|
| `0.5.x-1.20` | `1.20.x` | Java 11+ | Java 11 |
================================================================================
DOCUMENT: Installation (/docs/setup/setup-a-project)
================================================================================
# Installation
Learn how to configure your Maven `pom.xml` with the Flinkboot Bill of Materials (BOM), declare required modules, and package a production-ready Fat JAR.
---
## 1. Import the BOM
Add the Flinkboot Bill of Materials (BOM) to your project's ``. The BOM centrally aligns versions and configures Apache Flink runtime libraries with `provided` scope automatically.
```xml
io.github.sekelenao
flinkboot
0.5.0-1.20
pom
import
```
---
## 2. Declare dependencies
Add the Flinkboot modules you need to your `` section without specifying ``:
```xml
io.github.sekelenao
flinkboot-core
io.github.sekelenao
flinkboot-kafka
io.github.sekelenao
flinkboot-test
org.apache.flink
flink-streaming-java
org.apache.flink
flink-clients
org.slf4j
slf4j-api
```
---
## 3. Package the fat JAR
When deploying to an Apache Flink cluster, package your application with `maven-shade-plugin`.
Because Flink ships its own internal version of Jackson, relocate Jackson classes (`com.fasterxml`) to avoid classpath collisions at runtime:
```xml
org.apache.maven.plugins
maven-shade-plugin
3.6.0
package
shade
false
false
com.fasterxml
io.github.sekelenao.flinkboot.shaded.fasterxml
*:*
META-INF/*.SF
META-INF/*.DSA
META-INF/*.RSA
module-info.class
META-INF/versions/**
```
================================================================================
DOCUMENT: Collecting sink for integration testing (/docs/testing/collecting-sink-for-integration-testing)
================================================================================
# Collecting sink for integration testing
Testing streaming topologies often requires asserting on the actual data elements produced by a pipeline. Flinkboot provides the thread-safe `CollectingSink` utility in `flinkboot-test` to collect stream elements during integration tests.
---
## 1. Overview
`CollectingSink` is an in-memory, thread-safe Flink sink designed specifically for testing:
* `new CollectingSink()`: Creates a new sink instance.
* `sink.elements()`: Returns an immutable snapshot list of all collected elements.
* `sink.clear()`: Empties the internal buffer for subsequent test runs.
* `AutoCloseable`: Implements `AutoCloseable` for automatic buffer cleanup via `try-with-resources`.
The sink safely collects elements emitted across parallel sink subtasks during local Flink MiniCluster test execution.
---
## 2. Usage in integration tests
Attach `CollectingSink` to your stream with `.sinkTo(sink)`, execute the pipeline, and assert on the collected elements:
```java
import io.github.sekelenao.flinkboot.test.api.sink.CollectingSink;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.junit.jupiter.api.Test;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertAll;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class StreamPipelineTest {
@Test
void shouldProcessAndCollectEvents() throws Exception {
var env = StreamExecutionEnvironment.getExecutionEnvironment();
try (var sink = new CollectingSink()) {
env.fromData("event-1", "event-2")
.sinkTo(sink)
.setParallelism(2);
env.execute();
List elements = sink.elements();
assertAll(
() -> assertEquals(2, elements.size()),
() -> assertTrue(elements.containsAll(List.of("event-1", "event-2")))
);
}
}
}
```
Because Flink does not guarantee global ordering across parallel subtasks, assertions should verify element presence and count rather than relying on strict index ordering.
---
## 3. Buffer management & snapshots
* **Immutable snapshots**: Calling `sink.elements()` creates an immutable list copy. Further emissions or buffer clear operations do not mutate previously obtained snapshots.
* **Automatic resource cleanup**: Using `try-with-resources` automatically reclaims the in-memory sink storage once the test block completes.
* **Manual reuse with `clear()`**: If executing multiple pipelines sequentially in the same test method, call `sink.clear()` to reset the buffer between runs.
================================================================================
DOCUMENT: Loading configuration in tests (/docs/testing/load-configurations-in-tests)
================================================================================
# Loading configuration in tests
Flinkboot allows you to load, merge, and validate YAML configurations directly within your JUnit 5 tests using `Flinkboot.initialize(...)`.
---
## 1. Overview
When writing unit or integration tests for your Flink applications, you often need to load and validate your application configuration objects (DTOs) without starting a full command-line application:
* **Production parity**: Tests execute through the exact same parsing and validation engine used in production.
* **Varargs simplicity**: Call `Flinkboot.initialize()` with zero arguments for default classpath configuration, or pass command-line options inline.
* **Full option support**: Test custom flags (`--dry-run`), CLI parameters (`-threshold 100`), or Flinkboot runtime options alongside your configuration files.
* **Explicit scheme prefixes**: Unified support for `classpath:`, `resource:`, and `file:` schemes via Flinkboot's `Resource` API.
---
## 2. Usage examples
### Loading default classpath configuration
Calling `Flinkboot.initialize()` without arguments automatically loads `src/test/resources/job-configuration.yaml`:
```java
import io.github.sekelenao.flinkboot.core.api.Flinkboot;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertNotNull;
class ApplicationConfigTest {
@Test
void shouldLoadDefaultConfiguration() throws Exception {
MyApplicationConfig config = Flinkboot.initialize()
.configuration(MyApplicationConfig.class);
assertNotNull(config);
}
}
```
---
### Loading a specific configuration file
Pass `-flinkboot-configurations` with an explicit resource scheme prefix:
```java
@Test
void shouldLoadSpecificConfiguration() throws Exception {
MyApplicationConfig config = Flinkboot.initialize(
"-flinkboot-configurations", "classpath:job-test.yaml"
).configuration(MyApplicationConfig.class);
assertNotNull(config);
}
```
---
### Loading multiple merged configurations
You can pass multiple comma-separated configuration paths to test profile overrides:
```java
@Test
void shouldLoadMultipleConfigurations() throws Exception {
MyApplicationConfig config = Flinkboot.initialize(
"-flinkboot-configurations",
"classpath:config-base.yaml,classpath:config-test-env.yaml"
).configuration(MyApplicationConfig.class);
assertNotNull(config);
}
```
---
### Loading files from the file system
Target external configuration files outside the classpath using the `file:` prefix:
```java
@Test
void shouldLoadFromFileSystem() throws Exception {
MyApplicationConfig config = Flinkboot.initialize(
"-flinkboot-configurations", "file:/etc/flinkboot/my-config.yaml"
).configuration(MyApplicationConfig.class);
assertNotNull(config);
}
```
---
### Testing with flags and CLI parameters
You can pass custom flags and parameters inline:
```java
@Test
void shouldLoadConfigurationWithOptions() throws Exception {
Flinkboot boot = Flinkboot.initialize(
"-flinkboot-configurations", "classpath:job-test.yaml",
"--dry-run",
"-custom-param", "custom-value"
);
MyApplicationConfig config = boot.configuration(MyApplicationConfig.class);
assertAll(
() -> assertNotNull(config),
() -> assertTrue(boot.flag("dry-run")),
() -> assertEquals("custom-value", boot.parameter("custom-param").orElseThrow())
);
}
```
---
## 3. Scheme prefixes
Each path passed to `-flinkboot-configurations` must specify a valid resource scheme:
| Scheme prefix | Target location | Example |
|:--------------|:------------------------------------------------|:-------------------------------|
| `classpath:` | Classpath resources (e.g. `src/test/resources`) | `"classpath:job-test.yaml"` |
| `resource:` | Alias for classpath resources | `"resource:job-test.yaml"` |
| `file:` | Absolute or relative file system paths | `"file:/tmp/test-config.yaml"` |
Omitting the prefix will throw an `UnrecognizedResourceException`.
================================================================================
DOCUMENT: POJO compliance (/docs/testing/pojo-compliance)
================================================================================
# POJO compliance
Flinkboot provides test assertions to ensure your data classes strictly comply with Apache Flink's POJO serialization requirements and prevent runtime Kryo fallback.
---
## 1. Flink POJO requirements
Apache Flink uses an optimized serializer (`PojoSerializer`) for state and stream data exchange. When Flink recognizes a class as a valid POJO, it can perform direct field access and serialize elements with minimal overhead.
If a class does not satisfy Flink POJO criteria:
* Flink falls back to Kryo serialization, which is slower, less space-efficient, and risky for state schema evolution.
* Keyed operations on nested fields (e.g. `keyBy("fieldName")`) will fail.
A class must satisfy the following criteria:
1. The class must be **public** and standalone (or a `public static` nested class).
2. It must have a **public zero-argument constructor**.
3. All fields must be either **public** (non-final) or have **public getter and setter** methods following JavaBean conventions. Both styles are fully supported.
4. **No Kryo fallback fields**: No field or nested field may resolve to `GenericTypeInfo`. `FlinkbootAssertions.assertThat(MyClass.class).isPojo()` inspects fields recursively to guarantee pure native Flink serialization.
---
## 2. Usage in JUnit 5
Use `FlinkbootAssertions.assertThat(Class> type).isPojo()` to verify your model classes:
```java
import org.junit.jupiter.api.Test;
import static io.github.sekelenao.flinkboot.test.api.assertion.FlinkbootAssertions.assertThat;
class ModelComplianceTest {
@Test
void shouldComplyWithPojoRules() {
assertThat(MyModel.class).isPojo();
}
}
```
If the class violates any of Flink's requirements or contains fields falling back to Kryo serialization, the assertion fails immediately with a descriptive error message indicating the exact path of the invalid field.
---
## 3. Asserting generic types
A `Class` literal erases generic parameters, so `assertThat(Map.class)` cannot tell Flink which key and
value types are involved. Use a `TypeHint` to retain them:
```java
import org.apache.flink.api.common.typeinfo.TypeHint;
import java.util.List;
import java.util.Map;
import static io.github.sekelenao.flinkboot.test.api.assertion.FlinkbootAssertions.assertThat;
class ModelComplianceTest {
@Test
void shouldComplyWithGenericTypes() {
assertThat(new TypeHint