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

使用FlinkKafkaConsumer010与Kafka 1.1,如何在Flink程序中获取偏移量延迟信息?

Getting Kafka Consumer Offsets and Offset Lag with FlinkKafkaConsumer010 (Kafka 1.1)

Hey there! Let's break down how to get both individual record offsets and offset lag in your Flink application using FlinkKafkaConsumer010 for Kafka 1.1.

You have two main ways to access offset information depending on your use case:

a. Get Offset for Individual Records

If you need the offset of every message you consume, configure your FlinkKafkaConsumer010 to return raw ConsumerRecord objects instead of just the message value. This lets you directly pull the offset from each record.

Here's a working example:

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer010;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Properties;

public class OffsetExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("bootstrap.servers", "your-kafka-brokers");
        kafkaProps.setProperty("group.id", "flink-consumer-group");
        kafkaProps.setProperty("key.deserializer", StringDeserializer.class.getName());
        kafkaProps.setProperty("value.deserializer", StringDeserializer.class.getName());

        // Configure consumer to return full ConsumerRecord
        FlinkKafkaConsumer010<ConsumerRecord<String, String>> consumer = 
            new FlinkKafkaConsumer010<>(
                "your-topic",
                (record) -> new ConsumerRecord<>(
                    record.topic(), record.partition(), record.offset(),
                    new String(record.key()), new String(record.value())
                ),
                kafkaProps
            );

        DataStream<ConsumerRecord<String, String>> kafkaStream = env.addSource(consumer);

        // Extract offset from each record
        kafkaStream.map(record -> 
            String.format("Message: %s | Offset: %d", record.value(), record.offset())
        ).print();

        env.execute("Flink Kafka Offset Example");
    }
}

b. Get Committed Offsets for the Consumer Group

If you want the globally committed offsets (the position Flink has acknowledged it's processed up to), use Kafka's consumer API within a Flink rich function to fetch offsets your consumer group has submitted to Kafka.

import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import java.util.Collections;

public class CommittedOffsetReader extends RichMapFunction<ConsumerRecord<String,String>, String> {
    private transient KafkaConsumer<String, String> kafkaConsumer;
    private Properties kafkaProps;
    private String topic;

    public CommittedOffsetReader(Properties kafkaProps, String topic) {
        this.kafkaProps = kafkaProps;
        this.topic = topic;
    }

    @Override
    public void open(Configuration parameters) {
        // Initialize Kafka consumer to fetch committed offsets
        kafkaProps.setProperty("group.id", "flink-consumer-group");
        kafkaConsumer = new KafkaConsumer<>(kafkaProps);
        kafkaConsumer.subscribe(Collections.singletonList(topic));
    }

    @Override
    public String map(ConsumerRecord<String, String> record) {
        TopicPartition tp = new TopicPartition(record.topic(), record.partition());
        OffsetAndMetadata committed = kafkaConsumer.committed(tp);
        
        return String.format(
            "Partition %d | Current Record Offset: %d | Committed Offset: %d",
            record.partition(), record.offset(), 
            committed != null ? committed.offset() : 0
        );
    }

    @Override
    public void close() {
        if (kafkaConsumer != null) {
            kafkaConsumer.close();
        }
    }
}

2. Calculating Offset Lag

Offset lag is the gap between the latest available offset in a Kafka partition and the offset your consumer group has committed. Here's how to compute it in your Flink code:

import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import java.util.Collections;
import java.util.Map;

public class OffsetLagCalculator extends RichMapFunction<ConsumerRecord<String,String>, String> {
    private transient KafkaConsumer<String, String> kafkaConsumer;
    private Properties kafkaProps;
    private String topic;

    public OffsetLagCalculator(Properties kafkaProps, String topic) {
        this.kafkaProps = kafkaProps;
        this.topic = topic;
    }

    @Override
    public void open(Configuration parameters) {
        kafkaProps.setProperty("group.id", "flink-consumer-group");
        kafkaConsumer = new KafkaConsumer<>(kafkaProps);
        kafkaConsumer.subscribe(Collections.singletonList(topic));
    }

    @Override
    public String map(ConsumerRecord<String, String> record) {
        TopicPartition tp = new TopicPartition(record.topic(), record.partition());
        
        // Get latest offset in the partition
        Map<TopicPartition, Long> endOffsets = kafkaConsumer.endOffsets(Collections.singletonList(tp));
        Long latestOffset = endOffsets.get(tp);
        
        // Get committed offset for the consumer group
        OffsetAndMetadata committedOffset = kafkaConsumer.committed(tp);
        Long currentOffset = committedOffset != null ? committed.offset() : 0L;
        
        // Calculate lag
        long lag = latestOffset - currentOffset;

        return String.format(
            "Topic: %s | Partition: %d | Lag: %d",
            record.topic(), record.partition(), lag
        );
    }

    @Override
    public void close() {
        if (kafkaConsumer != null) {
            kafkaConsumer.close();
        }
    }
}

Key Notes to Remember

  • Performance: Calculating lag for every record adds overhead. Instead, use a ProcessFunction with timers to compute lag periodically (e.g., every 30 seconds) instead of on each message.
  • Initialization: If your consumer group is new, committedOffset might be null—handle this by defaulting to 0 or matching your consumer's auto-offset-reset configuration.
  • Compatibility: Ensure your Flink version is compatible with FlinkKafkaConsumer010 and Kafka 1.1 (Flink 1.7.x is a safe match for this combination).
  • Metrics Alternative: FlinkKafkaConsumer010 automatically exposes lag metrics through Flink's metric system. You can access these via a custom MetricReporter or Flink's REST API if you don't need lag directly in your application code.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:15:50