Skip to main content

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:

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

2. Kafka source​

YAML configuration​

You can configure subscriptions using either an explicit topic list or a regex pattern:

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:

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:

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 KeyTypeRequiredValidationDescription
nameStringYes@NotBlankOperator name registered in the Flink DAG graph.
bootstrap-serversList<String>Yes@NotEmpty, items @NotBlankKafka bootstrap broker hosts and ports.
group-idStringYes@NotBlankConsumer group ID.
topicsList<String>Conditionalitems @NotBlankExplicit topic subscriptions (mutually exclusive with topic-pattern).
topic-patternStringConditionalValid regexTopic subscription regex pattern (mutually exclusive with topics).
starting-offsetsEnumYes@NotNullStartup strategy: EARLIEST, LATEST, COMMITTED, TIMESTAMP, OFFSETS.
starting-offsets-timestampLongConditional@PositiveOrZeroTimestamp in milliseconds (mandatory if starting-offsets: TIMESTAMP, forbidden otherwise).
starting-offsets-partition-offsetsListConditional@Valid itemsList of partition starting offsets (mandatory if starting-offsets: OFFSETS, forbidden otherwise).
propertiesMap<String, String>NoFree-form mapAdditional Kafka consumer tuning properties (e.g. session.timeout.ms).

Offset strategies (starting-offsets)​

StrategyRequired Complementary KeysForbidden Complementary KeysDescription
EARLIESTNonestarting-offsets-timestamp, starting-offsets-partition-offsetsStart from earliest available log offsets.
LATESTNonestarting-offsets-timestamp, starting-offsets-partition-offsetsStart from latest log offsets.
COMMITTEDNonestarting-offsets-timestamp, starting-offsets-partition-offsetsStart from consumer group committed offsets.
TIMESTAMPstarting-offsets-timestamp (Long)starting-offsets-partition-offsetsPosition based on record epoch millisecond timestamps.
OFFSETSstarting-offsets-partition-offsets (List)starting-offsets-timestampExplicit starting offsets per topic partition.

3. Kafka sink​

YAML configuration​

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 KeyTypeRequiredValidationDescription
nameStringYes@NotBlankLogical identifier for the sink operator.
bootstrap-serversList<String>Yes@NotEmpty, items @NotBlankKafka broker endpoints.
topicStringYes@NotBlankTarget Kafka topic for emitted events.
delivery-guaranteeEnumYesNONE, AT_LEAST_ONCE, EXACTLY_ONCEDelivery semantic guarantee.
transactional-id-prefixStringConditionalStringTransactional prefix. Mandatory if delivery-guarantee is EXACTLY_ONCE, prohibited otherwise.
propertiesMap<String, String>NoNon-blank keys/valuesCustom Kafka producer client settings.

4. Pipeline integration​

Combine KafkaSourceProperties and KafkaSinkProperties in your configuration model and instantiate them via KafkaSourceFactory and KafkaSinkFactory:

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:

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<String> source = KafkaSourceFactory.supplyFor(
config.kafkaSource(),
KafkaRecordDeserializationSchema.valueOnly(new SimpleStringSchema())
);

// 2. Build Kafka Sink
KafkaSink<String> 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(...).