{"record":{"id":"778e3e9860db4f03","repo":"apache/beam","slug":"the-version-kafka-clients-does-not-support-record-headers","errorCode":null,"errorMessage":"The version kafka-clients does not support record headers, please use version 0.11.0.0 or newer","messagePattern":"The version kafka-clients does not support record headers, please use version 0\\.11\\.0\\.0 or newer","errorType":"exception","errorClass":"RuntimeException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaRecord.java","lineNumber":88,"sourceCode":"    this.kv = kv;\n  }\n\n  public String getTopic() {\n    return topic;\n  }\n\n  public int getPartition() {\n    return partition;\n  }\n\n  public long getOffset() {\n    return offset;\n  }\n\n  @Pure\n  public @Nullable Headers getHeaders() {\n    if (!ConsumerSpEL.hasHeaders()) {\n      throw new RuntimeException(\n          \"The version kafka-clients does not support record headers, \"\n              + \"please use version 0.11.0.0 or newer\");\n    }\n    return headers;\n  }\n\n  public KV<K, V> getKV() {\n    return kv;\n  }\n\n  public long getTimestamp() {\n    return timestamp;\n  }\n\n  public KafkaTimestampType getTimestampType() {\n    return timestampType;\n  }\n","sourceCodeStart":70,"sourceCodeEnd":106,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaRecord.java#L70-L106","documentation":"KafkaRecord.getHeaders() requires the kafka-clients Consumer to support record headers, an API added in Kafka 0.11.0.0. ConsumerSpEL.hasHeaders() checks this reflectively; when the classpath has an older kafka-clients jar, a RuntimeException is thrown.","triggerScenarios":"Calling getHeaders(), recordHeaders(), headers(), keys(), or building a structuralValue of a KafkaRecord while kafka-clients < 0.11.0.0 is on the classpath.","commonSituations":"Legacy builds pinning kafka-clients 0.9/0.10, transitive dependency conflicts downgrading kafka-clients, or brokers/clients from very old Kafka deployments.","solutions":["Upgrade the kafka-clients dependency to 0.11.0.0 or newer","Run `mvn dependency:tree` (or Gradle equivalent) to find and remove transitive pins that force an old kafka-clients version","If headers are not needed, avoid header-related APIs (getHeaders/keys) in the pipeline"],"exampleFix":"// before (pom.xml)\n<dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId><version>0.10.2.1</version></dependency>\n// after\n<dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId><version>3.6.1</version></dependency>","handlingStrategy":"validation","validationCode":"// check kafka-clients version at startup\nString version = org.apache.kafka.clients.consumer.ConsumerConfig.class.getPackage().getImplementationVersion();\nif (version == null || version.compareTo(\"0.11.0.0\") < 0) throw new IllegalStateException(\"kafka-clients >= 0.11.0.0 required\");","typeGuard":"boolean headersSupported() { return org.apache.beam.sdk.io.kafka.ConsumerSpEL.hasHeaders(); }","tryCatchPattern":"try {\n  Headers h = record.getHeaders();\n} catch (RuntimeException e) {\n  // kafka-clients too old; upgrade dependency\n}","preventionTips":["Pin kafka-clients to a modern version (>= 0.11.0.0)","Run dependency:tree to detect version conflicts/downgrades before release","Avoid header APIs when supporting legacy client deployments"],"tags":["java","kafka","dependency-version","headers"],"backgroundTag":"missing-dependency","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}