You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Java版Apache Flink处理Kafka多结构数据并路由?

处理Kafka多格式消息并路由到不同主题的Flink实现

问题1:能否读取多格式消息并路由到不同主题?

可以实现。核心思路是避免在数据源阶段绑定单一POJO反序列化器,先以原始字符串形式读取所有消息,再通过解析JSON结构判断消息类型,之后针对性地转换为对应数据对象,最后将不同类型的消息路由到指定Kafka主题。


问题2:Java版Flink代码实现

1. 定义对应两种消息结构的POJO

分别创建与两种JSON结构匹配的实体类(需提供无参构造器,用于Flink反序列化):

// 对应 {name: "abc", age: 20} 格式
public class UserBasic {
    private String name;
    private int age;

    public UserBasic() {}

    // Getter、Setter方法
    public String getName() { return name; }
    public void setName(String name) { this.name = name; }
    public int getAge() { return age; }
    public void setAge(int age) { this.age = age; }
}

// 对应 {pin: 111, number: 999999, address: "some place"} 格式
public class UserDetails {
    private int pin;
    private long number;
    private String address;

    public UserDetails() {}

    // Getter、Setter方法
    public int getPin() { return pin; }
    public void setPin(int pin) { this.pin = pin; }
    public long getNumber() { return number; }
    public void setNumber(long number) { this.number = number; }
    public String getAddress() { return address; }
    public void setAddress(String address) { this.address = address; }
}

2. 主流程实现代码

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.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.util.Collector;
import org.apache.flink.util.OutputTag;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.flink.api.common.serialization.SimpleStringSchema;

import java.util.Properties;

public class MultiFormatKafkaRouter {
    // 定义侧输出流标签,用于区分不同类型的消息
    private static final OutputTag<UserBasic> USER_BASIC_TAG = new OutputTag<UserBasic>("user-basic") {};
    private static final OutputTag<UserDetails> USER_DETAILS_TAG = new OutputTag<UserDetails>("user-details") {};
    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 配置Kafka源参数
        String bootstrapServers = "your-kafka-bootstrap-servers";
        String inputTopic = "topic_a";
        String groupId = "all-events-group-id";

        Properties sourceProps = new Properties();
        // 可添加额外消费者配置,如auto.offset.reset等

        // 以原始字符串形式读取Kafka消息
        KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
                .setBootstrapServers(bootstrapServers)
                .setTopics(inputTopic)
                .setProperties(sourceProps)
                .setGroupId(groupId)
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setValueOnlyDeserializer(new SimpleStringSchema())
                .build();

        DataStream<String> rawMessageStream = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Multi-Format Kafka Source");

        // 处理消息:区分类型并发送到对应侧输出流
        SingleOutputStreamOperator<Void> processedStream = rawMessageStream.process(new ProcessFunction<String, Void>() {
            @Override
            public void processElement(String rawMsg, Context ctx, Collector<Void> out) throws Exception {
                JsonNode jsonNode = OBJECT_MAPPER.readTree(rawMsg);

                // 通过关键字段判断消息类型
                if (jsonNode.has("name") && jsonNode.has("age")) {
                    UserBasic userBasic = OBJECT_MAPPER.treeToValue(jsonNode, UserBasic.class);
                    ctx.output(USER_BASIC_TAG, userBasic);
                } else if (jsonNode.has("pin") && jsonNode.has("number") && jsonNode.has("address")) {
                    UserDetails userDetails = OBJECT_MAPPER.treeToValue(jsonNode, UserDetails.class);
                    ctx.output(USER_DETAILS_TAG, userDetails);
                } else {
                    // 处理不符合格式的消息,可选择丢弃或发送到死信队列
                    System.err.println("Unrecognized message format: " + rawMsg);
                }
            }
        });

        // 从侧输出流获取对应类型的数据流
        DataStream<UserBasic> userBasicStream = processedStream.getSideOutput(USER_BASIC_TAG);
        DataStream<UserDetails> userDetailsStream = processedStream.getSideOutput(USER_DETAILS_TAG);

        // 配置Kafka生产者参数
        Properties sinkProps = new Properties();
        // 添加生产者配置,如acks、retries等

        // 创建发送到user_basic主题的Sink
        KafkaSink<UserBasic> userBasicSink = KafkaSink.<UserBasic>builder()
                .setBootstrapServers(bootstrapServers)
                .setKafkaProducerConfig(sinkProps)
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.builder()
                                .setTopic("user_basic")
                                .setValueSerializationSchema((elem, ctx) -> OBJECT_MAPPER.writeValueAsBytes(elem))
                                .build()
                )
                .build();

        // 创建发送到user_details主题的Sink
        KafkaSink<UserDetails> userDetailsSink = KafkaSink.<UserDetails>builder()
                .setBootstrapServers(bootstrapServers)
                .setKafkaProducerConfig(sinkProps)
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.builder()
                                .setTopic("user_details")
                                .setValueSerializationSchema((elem, ctx) -> OBJECT_MAPPER.writeValueAsBytes(elem))
                                .build()
                )
                .build();

        // 将数据流输出到对应Kafka主题
        userBasicStream.sinkTo(userBasicSink);
        userDetailsStream.sinkTo(userDetailsSink);

        env.execute("Kafka Multi-Format Message Router");
    }
}

关键依赖说明

需在项目中添加以下Maven依赖:

<!-- Flink Kafka连接器 -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka</artifactId>
    <version>${flink.version}</version>
</dependency>
<!-- Jackson JSON处理库 -->
<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-databind</artifactId>
    <version>2.15.2</version>
</dependency>

内容的提问来源于stack exchange,提问作者Abhinav

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 10:10:31