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

如何使用Apache Flink KafkaSource读取带键的Kafka记录?

要读取带键的Kafka记录并获取key、value、timestamp等信息,你需要修改KafkaSource的反序列化配置,不要用valueOnly的方式,而是改用能解析完整Kafka记录的反序列化器。以下是调整后的完整示例代码:

import org.apache.flink.connector.kafka.source.KafkaRecord;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.MapFunction;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.kafka.common.serialization.StringDeserializer;

public class KafkaKeyedRecordReader {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        String ip = "localhost:9092";

        // 构建能读取完整Kafka记录的Source
        KafkaSource<KafkaRecord<String, String>> source = KafkaSource.<KafkaRecord<String, String>>builder()
                .setBootstrapServers(ip)
                .setTopics("test3")
                .setGroupId("1")
                .setStartingOffsets(OffsetsInitializer.earliest())
                .setDeserializer(
                        KafkaRecordDeserializationSchema.builder()
                                .setKeyDeserializer(new StringDeserializer())
                                .setValueDeserializer(new StringDeserializer())
                                .build()
                )
                .build();

        // 获取包含完整Kafka记录的DataStream
        DataStream<KafkaRecord<String, String>> stream = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");

        // 解析并打印记录的key、value、timestamp等信息
        stream.map((MapFunction<KafkaRecord<String, String>, String>) record -> {
            String key = record.getKey();
            String value = record.getValue();
            long timestamp = record.getTimestamp();
            // 还可以获取offset、partition等其他元数据
            int partition = record.getPartition();
            long offset = record.getOffset();
            return String.format("Kafka Record - Key: %s, Value: %s, Timestamp: %d, Partition: %d, Offset: %d",
                    key, value, timestamp, partition, offset);
        }).print();

        env.execute("Read Keyed Kafka Records");
    }
}

关键改动说明

  • 泛型调整:将KafkaSource和DataStream的泛型从String改为KafkaRecord<String, String>,确保能接收完整的Kafka记录对象
  • 反序列化器替换:用KafkaRecordDeserializationSchema.builder()替代原来的valueOnly,分别指定key和value的反序列化器(这里用StringDeserializer匹配生产端的字符串格式)
  • 记录信息提取:在map算子中直接调用record.getKey()、record.getValue()、record.getTimestamp()等方法,还可以获取分区、偏移量等更多元数据

内容的提问来源于stack exchange,提问作者4 3 2

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 04:01:28