Flink 连接 Kafka:数据源接入与数据写入的完整代码示例
·
Flink 连接 Kafka:完整代码示例
以下是一个完整的 Flink 作业示例,包含从 Kafka 读取数据(Source)和处理后写入 Kafka(Sink)的全流程,使用 Java API 实现:
1. Maven 依赖配置
<dependencies>
<!-- Flink 核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.16.1</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.16.1</version>
</dependency>
<!-- Kafka 连接器 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.16.1</version>
</dependency>
<!-- JSON 处理 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-json</artifactId>
<version>1.16.1</version>
</dependency>
</dependencies>
2. 完整 Java 代码
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.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class KafkaFlinkIntegration {
public static void main(String[] args) throws Exception {
// 1. 创建执行环境
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 2. 配置 Kafka Source(数据源)
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka-broker1:9092,kafka-broker2:9092") // Kafka 集群地址
.setTopics("input-topic") // 输入主题
.setGroupId("flink-consumer-group") // 消费者组
.setStartingOffsets(OffsetsInitializer.earliest()) // 从最早偏移量开始
.setValueOnlyDeserializer(new SimpleStringSchema()) // 反序列化器
.build();
// 3. 创建数据流
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(), // 水印策略
"Kafka Source"
);
// 4. 数据处理(示例:转换为大写)
DataStream<String> processedStream = stream
.map(String::toUpperCase) // 数据处理逻辑
.name("uppercase-transform"); // 算子名称
// 5. 配置 Kafka Sink(数据输出)
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers("kafka-broker1:9092,kafka-broker2:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("output-topic") // 输出主题
.setValueSerializationSchema(new SimpleStringSchema())
.build()
)
.build();
// 6. 将处理结果写入 Kafka
processedStream.sinkTo(sink).name("kafka-sink");
// 7. 执行作业
env.execute("Flink-Kafka Integration Job");
}
}
3. 关键组件说明
-
Kafka Source 配置:
setBootstrapServers():Kafka 集群地址setTopics():订阅的输入主题setGroupId():消费者组 IDsetStartingOffsets():偏移量起始位置(支持earliest/latest/指定时间戳)
-
数据处理:
- 示例使用
map()进行简单的字符串大写转换 - 可替换为实际业务逻辑(如过滤、聚合、窗口计算等)
- 示例使用
-
Kafka Sink 配置:
setRecordSerializer():定义消息序列化方式setTopic():指定输出主题- 支持精确一次语义(需开启 Flink Checkpoint 和 Kafka 事务)
4. 运行准备
- 创建 Kafka 主题:
# 输入主题
kafka-topics.sh --create --topic input-topic \
--bootstrap-server kafka-broker:9092 \
--partitions 3 --replication-factor 2
# 输出主题
kafka-topics.sh --create --topic output-topic \
--bootstrap-server kafka-broker:9092 \
--partitions 3 --replication-factor 2
- 生产测试数据:
kafka-console-producer.sh --topic input-topic \
--bootstrap-server kafka-broker:9092
> test message 1
> hello flink
- 消费输出结果:
kafka-console-consumer.sh --topic output-topic \
--bootstrap-server kafka-broker:9092 \
--from-beginning
TEST MESSAGE 1
HELLO FLINK
5. 高级配置选项
-
精确一次语义:
env.enableCheckpointing(5000); // 每5秒做一次Checkpoint sink = KafkaSink.<String>builder() ... .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix("flink-tx-") .build(); -
JSON 数据处理:
// 使用 JSON 反序列化 KafkaSource<JsonNode> jsonSource = KafkaSource.builder() .setDeserializer(KafkaRecordDeserializationSchema.valueOnly( new JSONKeyValueDeserializationSchema(false))) ... -
动态主题选择:
.setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopicSelector((element) -> "topic-" + element.getPartition()) ...
此示例展示了 Flink 与 Kafka 集成的完整流程,可根据实际需求调整数据处理逻辑和配置参数。
更多推荐
所有评论(0)