使用FlinkKafkaConsumer010与Kafka 1.1,如何在Flink程序中获取偏移量延迟信息?
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.
1. Retrieving Kafka Consumer Offsets in Flink
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
ProcessFunctionwith timers to compute lag periodically (e.g., every 30 seconds) instead of on each message. - Initialization: If your consumer group is new,
committedOffsetmight benull—handle this by defaulting to 0 or matching your consumer'sauto-offset-resetconfiguration. - Compatibility: Ensure your Flink version is compatible with
FlinkKafkaConsumer010and Kafka 1.1 (Flink 1.7.x is a safe match for this combination). - Metrics Alternative:
FlinkKafkaConsumer010automatically exposes lag metrics through Flink's metric system. You can access these via a customMetricReporteror Flink's REST API if you don't need lag directly in your application code.
内容的提问来源于stack exchange,提问作者william

