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

如何在Apache Flink DataStream API中访问Kafka元数据?

你可以通过两种方式在DataStream API中获取Kafka的分区、偏移量、时间戳等元数据,无需依赖Table API:

Flink 1.17及以上版本的KafkaSource引入了KafkaRecord类型,直接封装了Kafka记录的元数据和内容。你可以通过以下方式获取:

KafkaSource<String> kafkaSource = KafkaSource.<String>builder()
    .setBootstrapServers("localhost:9092")
    .setTopics("your-topic")
    .setGroupId("flink-consumer-group")
    .setDeserializer(KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class))
    .build();

// 从KafkaSource获取包含元数据的KafkaRecord流
DataStream<KafkaRecord<String>> recordStream = env.fromSource(
    kafkaSource,
    WatermarkStrategy.noWatermarks(),
    "Kafka Source With Metadata"
);

// 提取元数据和业务数据
recordStream.map(kafkaRecord -> {
    String value = kafkaRecord.value();
    int partition = kafkaRecord.partition();
    long offset = kafkaRecord.offset();
    long timestamp = kafkaRecord.timestamp();
    
    return String.format("业务数据:%s,分区:%d,偏移量:%d,时间戳:%d", value, partition, offset, timestamp);
}).print();

在旧版本中,你需要自定义反序列化Schema,直接从ConsumerRecord中提取元数据,并将业务数据和元数据封装成自定义对象返回:

步骤1:定义元数据封装类

public class KafkaRecordMeta {
    private int partition;
    private long offset;
    private long timestamp;
    private String value;

    // 构造方法、getter/setter、toString
    public KafkaRecordMeta(String value, int partition, long offset, long timestamp) {
        this.value = value;
        this.partition = partition;
        this.offset = offset;
        this.timestamp = timestamp;
    }

    // 省略getter、setter和toString方法
}

步骤2:自定义反序列化Schema

public class MetaAwareKafkaDeserializer implements KafkaDeserializationSchema<KafkaRecordMeta> {
    @Override
    public boolean isEndOfStream(KafkaRecordMeta nextElement) {
        return false;
    }

    @Override
    public KafkaRecordMeta deserialize(ConsumerRecord<byte[], byte[]> consumerRecord) throws Exception {
        String value = new String(consumerRecord.value(), StandardCharsets.UTF_8);
        return new KafkaRecordMeta(
            value,
            consumerRecord.partition(),
            consumerRecord.offset(),
            consumerRecord.timestamp()
        );
    }

    @Override
    public TypeInformation<KafkaRecordMeta> getProducedType() {
        return TypeInformation.of(KafkaRecordMeta.class);
    }
}

步骤3:使用自定义Schema创建数据源

Properties kafkaProps = new Properties();
kafkaProps.setProperty("bootstrap.servers", "localhost:9092");
kafkaProps.setProperty("group.id", "flink-old-version-group");

FlinkKafkaConsumer<KafkaRecordMeta> consumer = new FlinkKafkaConsumer<>(
    "your-topic",
    new MetaAwareKafkaDeserializer(),
    kafkaProps
);

DataStream<KafkaRecordMeta> stream = env.addSource(consumer);

// 处理包含元数据的流
stream.map(record -> {
    // 业务逻辑处理
    return record.toString();
}).print();

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 20:22:32