Flink 连接 Kafka:数据源接入与数据写入的完整代码示例
·
Flink 连接 Kafka:数据源接入与数据写入的完整代码示例
以下是使用 Apache Flink 连接 Kafka 的完整代码示例,包括数据源接入(从 Kafka 读取数据)和数据写入(将数据写入 Kafka)。代码基于 Flink 1.17.x 和 Kafka 客户端库,使用 Java 语言实现。示例包含详细注释,确保结构清晰。
前置条件
- 依赖管理:使用 Maven 添加以下依赖(在
pom.xml中):
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.4.0</version>
</dependency>
</dependencies>
- 环境准备:
- Kafka 集群运行中(例如
localhost:9092)。 - 创建输入主题(如
input-topic)和输出主题(如output-topic)。
- Kafka 集群运行中(例如
完整代码示例
以下代码实现了一个简单的 Flink 作业:
- 从 Kafka 读取数据(数据源接入)。
- 处理数据(例如,将字符串转为大写)。
- 将处理后的数据写入 Kafka(数据写入)。
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 FlinkKafkaIntegrationExample {
public static void main(String[] args) throws Exception {
// 1. 设置 Flink 执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Kafka 配置参数
String bootstrapServers = "localhost:9092"; // Kafka 服务器地址
String inputTopic = "input-topic"; // 输入主题名称
String outputTopic = "output-topic"; // 输出主题名称
String groupId = "flink-consumer-group"; // 消费者组 ID
// 2. 数据源接入:从 Kafka 读取数据
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers(bootstrapServers) // Kafka 服务器
.setTopics(inputTopic) // 订阅的主题
.setGroupId(groupId) // 消费者组
.setStartingOffsets(OffsetsInitializer.earliest()) // 从最早偏移量开始读取
.setValueOnlyDeserializer(new SimpleStringSchema()) // 反序列化器(字符串格式)
.build();
// 创建数据流,从 Kafka 源读取
DataStream<String> inputStream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(), // 水印策略(无状态处理)
"Kafka Source" // 算子名称
);
// 3. 数据处理:简单转换(示例:字符串转为大写)
DataStream<String> processedStream = inputStream
.map(String::toUpperCase) // 每个元素转为大写
.name("Data Processor"); // 算子名称
// 4. 数据写入:将结果写入 Kafka
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers(bootstrapServers) // Kafka 服务器
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic(outputTopic) // 输出主题
.setValueSerializationSchema(new SimpleStringSchema()) // 序列化器(字符串格式)
.build()
)
.build();
// 将处理后的数据流写入 Kafka Sink
processedStream.sinkTo(sink)
.name("Kafka Sink");
// 5. 执行 Flink 作业
env.execute("Flink Kafka Integration Job");
}
}
代码说明
-
数据源接入:
- 使用
KafkaSource构建器从 Kafka 读取数据。 setStartingOffsets(OffsetsInitializer.earliest()):从主题的起始偏移量读取,确保不遗漏数据。SimpleStringSchema:假设 Kafka 消息是字符串格式,可根据需求替换为其他序列化器(如 JSON)。
- 使用
-
数据处理:
- 示例中使用了
map算子将字符串转为大写,实际应用可替换为复杂逻辑(如过滤、聚合)。 - 算子名称(如
"Data Processor")便于 Flink Web UI 监控。
- 示例中使用了
-
数据写入:
- 使用
KafkaSink构建器将数据写入 Kafka。 setRecordSerializer配置序列化方式,确保输出格式匹配 Kafka 主题。
- 使用
运行步骤
- 编译与打包:使用 Maven 编译项目:
mvn clean package - 提交作业:将生成的 JAR 文件提交到 Flink 集群:
flink run -c FlinkKafkaIntegrationExample path/to/your-jar-file.jar - 测试数据流:
- 向
input-topic发送消息(如{"message": "hello"})。 - 从
output-topic消费消息,应看到大写结果(如HELLO)。
- 向
注意事项
- 错误处理:添加重试机制或死信队列处理异常数据。
- 性能优化:
- 调整 Flink 并行度(
env.setParallelism(4))。 - 启用检查点(Checkpointing)保证 Exactly-Once 语义。
- 调整 Flink 并行度(
- 安全配置:如果 Kafka 启用 SASL/SSL,需在
KafkaSource和KafkaSink中添加安全参数。 - 版本兼容:确保 Flink 和 Kafka 版本兼容,避免 API 不匹配问题。
此示例覆盖了 Flink 连接 Kafka 的核心场景,您可根据实际需求扩展数据处理逻辑或配置参数。
更多推荐


所有评论(0)