Flink 连接 Kafka:数据源接入与数据写入的完整代码示例

以下是使用 Apache Flink 连接 Kafka 的完整代码示例,包括数据源接入(从 Kafka 读取数据)和数据写入(将数据写入 Kafka)。代码基于 Flink 1.17.x 和 Kafka 客户端库,使用 Java 语言实现。示例包含详细注释,确保结构清晰。

前置条件
  1. 依赖管理:使用 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>

  1. 环境准备
    • Kafka 集群运行中(例如 localhost:9092)。
    • 创建输入主题(如 input-topic)和输出主题(如 output-topic)。
完整代码示例

以下代码实现了一个简单的 Flink 作业:

  1. 从 Kafka 读取数据(数据源接入)。
  2. 处理数据(例如,将字符串转为大写)。
  3. 将处理后的数据写入 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");
    }
}

代码说明
  1. 数据源接入

    • 使用 KafkaSource 构建器从 Kafka 读取数据。
    • setStartingOffsets(OffsetsInitializer.earliest()):从主题的起始偏移量读取,确保不遗漏数据。
    • SimpleStringSchema:假设 Kafka 消息是字符串格式,可根据需求替换为其他序列化器(如 JSON)。
  2. 数据处理

    • 示例中使用了 map 算子将字符串转为大写,实际应用可替换为复杂逻辑(如过滤、聚合)。
    • 算子名称(如 "Data Processor")便于 Flink Web UI 监控。
  3. 数据写入

    • 使用 KafkaSink 构建器将数据写入 Kafka。
    • setRecordSerializer 配置序列化方式,确保输出格式匹配 Kafka 主题。
运行步骤
  1. 编译与打包:使用 Maven 编译项目:
    mvn clean package
    

  2. 提交作业:将生成的 JAR 文件提交到 Flink 集群:
    flink run -c FlinkKafkaIntegrationExample path/to/your-jar-file.jar
    

  3. 测试数据流
    • input-topic 发送消息(如 {"message": "hello"})。
    • output-topic 消费消息,应看到大写结果(如 HELLO)。
注意事项
  1. 错误处理:添加重试机制或死信队列处理异常数据。
  2. 性能优化
    • 调整 Flink 并行度(env.setParallelism(4))。
    • 启用检查点(Checkpointing)保证 Exactly-Once 语义。
  3. 安全配置:如果 Kafka 启用 SASL/SSL,需在 KafkaSourceKafkaSink 中添加安全参数。
  4. 版本兼容:确保 Flink 和 Kafka 版本兼容,避免 API 不匹配问题。

此示例覆盖了 Flink 连接 Kafka 的核心场景,您可根据实际需求扩展数据处理逻辑或配置参数。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐