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:
strategy: EARLIEST
# Escape hatch: any native Kafka consumer property (SSL, SASL, timeouts)
properties:
security.protocol: "SASL_SSL"
sasl.mechanism: "SCRAM-SHA-512"
sasl.jaas.config: "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"${KAFKA_USER}\" password=\"${KAFKA_PASSWORD}\";"
ssl.truststore.location: "/var/private/ssl/kafka.truststore.jks"
ssl.truststore.password: "${KAFKA_TRUSTSTORE_PASSWORD}"
session.timeout.ms: "45000"

For bounded batch executions, specify boundedness: BOUNDED and configure stopping-offsets:

kafka-source:
name: "batch-orders-source"
bootstrap-servers:
- "localhost:9092"
group-id: "analytics-batch"
topic-pattern: "^analytics-.*$"
boundedness: BOUNDED
starting-offsets:
strategy: TIMESTAMP
timestamp: 1689717600000
stopping-offsets:
strategy: TIMESTAMP
timestamp: 1689721200000

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-offsetsKafkaOffsetPropertiesYes@NotNull @ValidOffset strategy used on startup.
boundednessEnumNoEnumUNBOUNDED (default) or BOUNDED.
stopping-offsetsKafkaOffsetPropertiesConditional@ValidStopping offset position. Mandatory when boundedness is BOUNDED.
propertiesMap<String, String>NoNon-blank keys/valuesEscape hatch passed to KafkaSourceBuilder.setProperties(...) (e.g. SSL, SASL, timeouts).

Offset positioning (starting-offsets / stopping-offsets)​

StrategyRequired ParametersProhibited ParametersDescription
EARLIESTNonetimestamp, partitionsStart from earliest available log offsets.
LATESTNonetimestamp, partitionsStart from latest log offsets.
COMMITTEDNonetimestamp, partitionsStart from consumer group committed offsets.
TIMESTAMPtimestamp (Long)partitionsPosition based on record epoch millisecond timestamps.
OFFSETSpartitions (Map<Integer, Long>)timestampExplicit mapping of partition indices to exact offsets.

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"
# Escape hatch: any native Kafka producer property (acks, compression, batching)
properties:
acks: "all"
compression.type: "zstd"
linger.ms: "20"

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/valuesEscape hatch passed to KafkaSinkBuilder.setKafkaProducerConfig(...) (e.g. acks, compression, batching).

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(...).