Skip to main content

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:

<dependencies>
<dependency>
<groupId>io.github.sekelenao</groupId>
<artifactId>flinkboot-fluss</artifactId>
</dependency>
</dependencies>

2. Fluss source​

YAML configuration​

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:

fluss-source:
name: "replay-source"
bootstrap-servers:
- "localhost:9123"
database: "analytics_db"
table: "user_events"
startup-mode: "TIMESTAMP"
startup-timestamp: 1700000000000

Configuration reference​

Property KeyTypeRequiredValidationDescription
nameStringYes@NotBlankUnique operator identifier in the Flink DAG graph.
bootstrap-serversList<String>Yes@NotEmpty, items @NotBlankAddresses of Fluss coordinators.
databaseStringYes@NotBlankTarget Fluss database.
tableStringYes@NotBlankTarget Fluss table.
startup-modeEnumYes@NotNullStartup strategy: EARLIEST, LATEST, FULL, TIMESTAMP.
startup-timestampLongConditional@PositiveOrZeroEpoch millisecond timestamp (mandatory if startup-mode: TIMESTAMP, forbidden otherwise).
propertiesMap<String, String>No@NotNull entriesAdditional Fluss client/scanner tuning options.

3. Fluss sink​

YAML configuration​

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 KeyTypeRequiredValidationDescription
nameStringYes@NotBlankUnique operator identifier in the Flink DAG graph.
bootstrap-serversList<String>Yes@NotEmpty, items @NotBlankAddresses of Fluss coordinators.
databaseStringYes@NotBlankTarget Fluss database.
tableStringYes@NotBlankTarget Fluss table.
propertiesMap<String, String>No@NotNull entriesAdditional Fluss client/writer tuning options.

4. Pipeline integration​

Combine FlussSourceProperties and FlussSinkProperties in your configuration model and initialize them with FlussSourceFactory and FlussSinkFactory:

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:

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<RowData> source = FlussSourceFactory.supplyFor(config.flussSource(), deserializer);

// 2. Build Fluss Sink
RowDataSerializationSchema serializer = ...; // Your serialization schema
FlussSink<RowData> sink = FlussSinkFactory.supplyFor(config.flussSink(), serializer);

// 3. Connect stream pipeline
DataStream<RowData> 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:

--add-opens=java.base/java.nio=ALL-UNNAMED

When running unit/integration tests with Maven Surefire:

<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<configuration>
<argLine>--add-opens=java.base/java.nio=ALL-UNNAMED</argLine>
</configuration>
</plugin>