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. 关键组件说明
  1. Kafka Source 配置

    • setBootstrapServers():Kafka 集群地址
    • setTopics():订阅的输入主题
    • setGroupId():消费者组 ID
    • setStartingOffsets():偏移量起始位置(支持 earliest/latest/指定时间戳)
  2. 数据处理

    • 示例使用 map() 进行简单的字符串大写转换
    • 可替换为实际业务逻辑(如过滤、聚合、窗口计算等)
  3. Kafka Sink 配置

    • setRecordSerializer():定义消息序列化方式
    • setTopic():指定输出主题
    • 支持精确一次语义(需开启 Flink Checkpoint 和 Kafka 事务)
4. 运行准备
  1. 创建 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

  1. 生产测试数据:
kafka-console-producer.sh --topic input-topic \
  --bootstrap-server kafka-broker:9092
> test message 1
> hello flink

  1. 消费输出结果:
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 集成的完整流程,可根据实际需求调整数据处理逻辑和配置参数。

更多推荐