获取提交偏移量遇SchemaException错误,求解决(目标:代码算Kafka主题总延迟)
coordinator_key SchemaException & Calculating Kafka Topic Total Lag Hey there, let's tackle your problem head-on. First, let's break down why that error is happening, then fix it, and finally get you the total lag you need.
Why You're Seeing the SchemaException
The root issue here is version incompatibility between your Kafka Broker (0.11.0.2) and Java client (1.0.0).
When you try to fetch committed offsets, the 1.0.0 client constructs a request that includes the coordinator_key field—this is a newer field introduced in later Kafka versions. But your 0.11.0.2 Broker's protocol schema doesn't recognize this field, and since it has no default value defined in the older schema, the Broker throws that error when trying to parse the request.
Solutions to Fix the Error
You have two solid options to resolve this:
Option 1: Align Client and Broker Versions (Most Reliable)
Downgrade your kafka-clients dependency to match your Broker version (0.11.0.2). This ensures full protocol compatibility, so there won't be any mismatched request fields. Update your pom.xml (or build.gradle) to use:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>0.11.0.2</version> </dependency>
Option 2: Force Client to Use Older Protocol Version
If you can't downgrade the client, configure it to use the older protocol version compatible with your Broker. Add these settings when initializing your KafkaConsumer:
Properties props = new Properties(); // ... your existing consumer configs (bootstrap.servers, group.id, etc.) props.put("api.version.request", "false"); props.put("broker.version.fallback", "0.11.0.2");
This tells the client not to request the latest API version and instead fall back to the 0.11.0.2 protocol, which won't include the problematic coordinator_key field.
Calculating Total Topic Lag
Once the error is fixed, you can compute the total lag for your topic with these steps:
Get all partitions for your target topic:
List<TopicPartition> partitions = consumer.partitionsFor("your-topic-name") .stream() .map(p -> new TopicPartition(p.topic(), p.partition())) .collect(Collectors.toList());Fetch the latest end offsets for each partition:
Map<TopicPartition, Long> endOffsets = consumer.endOffsets(partitions);Fetch the committed offsets for each partition:
Map<TopicPartition, OffsetAndMetadata> committedOffsets = consumer.committed(new HashSet<>(partitions));Calculate total lag by summing the difference per partition:
long totalLag = 0; for (TopicPartition partition : partitions) { Long endOffset = endOffsets.get(partition); OffsetAndMetadata committed = committedOffsets.get(partition); if (endOffset != null && committed != null) { totalLag += (endOffset - committed.offset()); } } System.out.println("Total topic lag: " + totalLag);
Quick Notes:
- Make sure your consumer is subscribed to the topic (or assigned partitions) before calling these methods.
- If you're checking lag for a specific consumer group, ensure the
group.idin your config matches that group.
内容的提问来源于stack exchange,提问作者Doron Levi

