温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

flink框架如何对接kafka

发布时间:2026-09-06 16:10:18 来源:亿速云 阅读:92 作者:小樊 栏目:软件技术

Apache Flink 中对接 Kafka 是最常见也是官方支持最好的场景之一。下面从依赖、连接方式、常见用法、参数调优、Exactly-Once几个方面给你一个完整说明。


一、引入 Kafka 连接器依赖

1️⃣ Flink 1.14+(推荐,使用 flink-connector-kafka

Maven

<dependency>
  <groupId>org.apache.flink</groupId>
  <artifactId>flink-connector-kafka</artifactId>
  <version>1.17.0</version>
</dependency>

⚠️ 注意:Kafka 连接器版本要和 Flink 版本匹配


二、Flink 消费 Kafka(Source)

1️⃣ 基础示例(DataStream API)

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

2️⃣ 常用参数说明

参数 说明
bootstrap.servers Kafka broker 地址
group.id 消费者组
startingOffsets earliest / latest / 指定 offset
deserializer 反序列化(String / JSON / Avro)

三、Flink 写入 Kafka(Sink)

1️⃣ Kafka Sink 示例

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);

四、Table API / SQL 方式(更常用)

1️⃣ 创建 Kafka 表

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'
);

2️⃣ 写入 Kafka

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;

五、Exactly-Once 保障

关键配置

env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

Kafka Sink:

.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)

✅ 前提:

  • 开启 Checkpoint
  • Kafka 0.11+
  • Flink 1.4+

六、常见问题

❓ Offset 不提交?

  • 没开 Checkpoint
  • group.id 重复

❓ 数据重复?

  • 使用 AT_LEAST_ONCE
  • Kafka 幂等 Producer 未开启

❓ 消费慢?

  • 增加并行度
  • 调整 fetch.max.wait.ms

七、适用场景总结

✅ 实时 ETL
✅ 日志处理
✅ 实时指标计算
✅ 事件驱动架构


如果你愿意,我可以:

  • 给你 完整 Demo 项目结构
  • Flink + Kafka + Checkpoint 底层原理
  • 或直接针对你的 Flink / Kafka 版本给配置

你现在是 DataStream 还是 SQL 方式?

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI