# 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>>() {}).isPojo(); } } ``` Every nested type is validated recursively: the assertion fails if the key type, the value type, or any nested field falls back to Kryo. --- ## 4. Asserting existing `TypeInformation` When a type description already exists (produced by a custom `TypeInfoFactory`, by `TypeInformation.of(...)`, or returned by an operator), pass it directly: ```java import org.apache.flink.api.common.typeinfo.TypeInformation; import static io.github.sekelenao.flinkboot.test.api.assertion.FlinkbootAssertions.assertThat; class ModelComplianceTest { @Test void shouldComplyWithTypeInformation() { TypeInformation typeInfo = TypeInformation.of(MyModel.class); assertThat(typeInfo).isPojo(); } } ``` This verifies directly that a custom factory produces a native Flink serializer instead of a Kryo fallback. ================================================================================ DOCUMENT: Serialization compliance (/docs/testing/serialization-compliance) ================================================================================ # Serialization compliance In Apache Flink streaming applications, user-defined functions and helper objects are distributed across remote TaskManagers using standard Java serialization. If any captured object or nested field does not implement `java.io.Serializable`, Flink fails at submission or initialization with a fatal `java.io.NotSerializableException`. Flinkboot provides testing assertions in `flinkboot-test` to verify that your classes and object graphs comply with standard Java serialization requirements. --- ## 1. How `isSerializable()` works The assertion `FlinkbootAssertions.assertThat(actual).isSerializable()` performs a complete in-memory serialization round-trip: 1. **Serialization**: Writes the target object into a byte buffer using `java.io.ObjectOutputStream`. 2. **Deserialization**: Reconstructs the object from the byte buffer using `java.io.ObjectInputStream`. This verifies that: * The target class implements `java.io.Serializable`. * Every nested object, collection element, and referenced type in the object graph is serializable. * Custom `writeObject` or `readObject` hooks execute without errors. If any element fails serialization, the assertion fails immediately with a descriptive error detailing the exact offending field. --- ## 2. Usage in JUnit 5 Pass any object instance directly to `assertThat(actual).isSerializable()`: ```java import org.junit.jupiter.api.Test; import static io.github.sekelenao.flinkboot.test.api.assertion.FlinkbootAssertions.assertThat; class SerializationComplianceTest { @Test void shouldBeSerializable() { assertThat(payload).isSerializable(); } } ``` ================================================================================ DOCUMENT: Loading resources (/docs/utilities/loading-resources) ================================================================================ # Loading resources Streaming applications frequently need to load side-inputs, reference datasets, encryption keys, and schema definitions across classpath and filesystem locations. Flinkboot provides the `Resource` API to unify and standardize these access patterns. --- ## 1. Unified resource API (`Resource`) The `Resource` contract allows reading content via an `InputStream` without managing manual URL or file descriptors: ```java import io.github.sekelenao.flinkboot.core.api.resource.Resource; import java.io.InputStream; import java.nio.charset.StandardCharsets; Resource rules = Resource.of("classpath:rules/fraud-rules.json"); try (InputStream in = rules.inputStream()) { String json = new String(in.readAllBytes(), StandardCharsets.UTF_8); } ``` --- ## 2. Supported scheme prefixes Flinkboot resolves paths based on scheme prefixes: | Scheme Prefix | Target Location | Example | |:---|:---|:---| | `classpath:` | JAR classpath resources | `Resource.of("classpath:schemas/event.avsc")` | | `resource:` | Alias for classpath resources | `Resource.of("resource:lookup.csv")` | | `file:` | Local or mounted filesystem | `Resource.of("file:/etc/secrets/tls.keystore")` | Prefix matching is case-insensitive (`file:`, `FILE:`, `classpath:`). --- ## 3. Distributed streaming lifecycle `Resource` instances are **not `Serializable`**. When accessing resources inside distributed Flink streaming operators (such as a `RichMapFunction` or `ProcessFunction`), do not serialize the `Resource` instance across TaskManagers. Instead, pass the path string to the operator constructor, declare the resource or resulting state as `transient`, and read the input stream inside the operator's `open(OpenContext)` lifecycle method. ================================================================================ DOCUMENT: Serializing JDK types (/docs/utilities/serializing-jdk-types) ================================================================================ # Serializing JDK types In Apache Flink, standard JDK date-time types and collections default to Kryo serialization, which is slower, less space-efficient, and risky for state schema evolution. Flinkboot provides built-in, optimized `TypeInfoFactory` classes to enable native Flink serialization for these types. --- ## 1. Available factories | Field Type | Flinkboot Factory Class | Serializer / Underlying Type | |:--------------------------|:-------------------------------|:------------------------------------| | `java.time.LocalDateTime` | `LocalDateTimeTypeInfoFactory` | `Types.LOCAL_DATE_TIME` | | `java.time.LocalDate` | `LocalDateTypeInfoFactory` | `Types.LOCAL_DATE` | | `java.time.LocalTime` | `LocalTimeTypeInfoFactory` | `Types.LOCAL_TIME` | | `java.time.Duration` | `DurationTypeInfoFactory` | Custom 12-byte `DurationSerializer` | | `java.util.List` | `ListTypeInfoFactory` | `Types.LIST(elementType)` | | `java.util.Map` | `MapTypeInfoFactory` | `Types.MAP(keyType, valueType)` | --- ## 2. Usage in POJO classes Annotate your POJO fields with Flink's `@TypeInfo` annotation: ```java import io.github.sekelenao.flinkboot.core.api.typing.collection.ListTypeInfoFactory; import io.github.sekelenao.flinkboot.core.api.typing.collection.MapTypeInfoFactory; import io.github.sekelenao.flinkboot.core.api.typing.time.DurationTypeInfoFactory; import io.github.sekelenao.flinkboot.core.api.typing.time.LocalDateTypeInfoFactory; import io.github.sekelenao.flinkboot.core.api.typing.time.LocalDateTimeTypeInfoFactory; import io.github.sekelenao.flinkboot.core.api.typing.time.LocalTimeTypeInfoFactory; import org.apache.flink.api.common.typeinfo.TypeInfo; import java.time.Duration; import java.time.LocalDate; import java.time.LocalDateTime; import java.time.LocalTime; import java.util.List; import java.util.Map; public class Example { @TypeInfo(LocalDateTimeTypeInfoFactory.class) public LocalDateTime eventTime; @TypeInfo(LocalDateTypeInfoFactory.class) public LocalDate eventDate; @TypeInfo(LocalTimeTypeInfoFactory.class) public LocalTime eventClock; @TypeInfo(DurationTypeInfoFactory.class) public Duration duration; @TypeInfo(ListTypeInfoFactory.class) public List tags; @TypeInfo(MapTypeInfoFactory.class) public Map metrics; } ``` --- ## 3. Edge cases and nuances ### A. Concrete implementations for `List` and `Map` Flink's built-in `ListSerializer` and `MapSerializer` instantiate `java.util.ArrayList` and `java.util.HashMap` upon deserialization. * If your POJO uses `java.util.List` or `java.util.Map`, it will be deserialized as an `ArrayList` or `HashMap`. * If your application relies on specific implementations (such as `java.util.TreeMap` for sorted keys or immutable collections), you must use a custom serializer. ### B. Nested JDK types inside collections In Java and Flink, the `@TypeInfo` annotation cannot be placed directly on generic type arguments (e.g. `List<@TypeInfo(...) Duration>`). To serialize `List` natively without Kryo, create a dedicated container factory: ```java public class DurationListTypeInfoFactory extends TypeInfoFactory> { @Override public TypeInformation> createTypeInfo(Type t, Map> genericParameters) { return Types.LIST(DurationTypeInfo.INSTANCE); } } ``` And annotate the field: ```java @TypeInfo(DurationListTypeInfoFactory.class) public List durations; ``` ### C. Custom domain classes in collections If you have a collection of custom DTOs (e.g. `List`), annotate `MyItem` at the class level: ```java @TypeInfo(MyItemTypeInfoFactory.class) public class MyItem { ... } ``` Flink's `TypeExtractor` will automatically find the class-level annotation when resolving `List`. ### D. Nullability Flink's native serializers for collections and time types properly support `null` values within POJO fields. ### E. Generic bounds vs wildcards in collections * **Bounded Class Generics (`Container`)**: When a concrete class argument is provided (e.g. `Container`), Flink's `TypeExtractor` resolves `ChildDto` natively as a POJO. * **Wildcards in Collections (`List`)**: Wildcard type arguments cannot be resolved into concrete type parameters by Flink's `TypeInfoFactory` and therefore fall back to Kryo serialization (`GenericTypeInfo`). Always declare collections with concrete type arguments (e.g. `List` instead of `List`).