如何使用Apache Flink KafkaSource读取带键的Kafka记录?
解决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
相关产品推荐
相关产品推荐

