在 Apache Flink 中对接 Kafka 是最常见也是官方支持最好的场景之一。下面从依赖、连接方式、常见用法、参数调优、Exactly-Once几个方面给你一个完整说明。
flink-connector-kafka)Maven
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.17.0</version>
</dependency>
⚠️ 注意:Kafka 连接器版本要和 Flink 版本匹配
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("localhost:9092")
.setTopics("test-topic")
.setGroupId("flink-group")
.setStartingOffsets(OffsetsInitializer.earliest())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"Kafka Source"
);
stream.print();
env.execute();
| 参数 | 说明 |
|---|---|
| bootstrap.servers | Kafka broker 地址 |
| group.id | 消费者组 |
| startingOffsets | earliest / latest / 指定 offset |
| deserializer | 反序列化(String / JSON / Avro) |
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers("localhost:9092")
.setRecordSerializer(
KafkaRecordSerializationSchema.builder()
.setTopic("output-topic")
.setValueSerializationSchema(new SimpleStringSchema())
.build()
)
.setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
.build();
stream.sinkTo(sink);
CREATE TABLE kafka_source (
id STRING,
name STRING,
ts TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'user-topic',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-sql',
'scan.startup.mode' = 'earliest-offset',
'format' = 'json'
);
CREATE TABLE kafka_sink (
id STRING,
cnt BIGINT
) WITH (
'connector' = 'kafka',
'topic' = 'result-topic',
'properties.bootstrap.servers' = 'localhost:9092',
'format' = 'json'
);
INSERT INTO kafka_sink
SELECT id, COUNT(*) FROM kafka_source GROUP BY id;
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
Kafka Sink:
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
✅ 前提:
AT_LEAST_ONCEfetch.max.wait.ms✅ 实时 ETL
✅ 日志处理
✅ 实时指标计算
✅ 事件驱动架构
如果你愿意,我可以:
你现在是 DataStream 还是 SQL 方式?
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。