{"record":{"id":"4987d93841890044","repo":"apache/beam","slug":"externalwithmetadata-transform-only-supports-values-of-type","errorCode":null,"errorMessage":"ExternalWithMetadata transform only supports values of type nullable(byte[])","messagePattern":"ExternalWithMetadata transform only supports values of type nullable\\(byte\\[\\]\\)","errorType":"exception","errorClass":"RuntimeException","httpStatus":null,"severity":"error","filePath":"sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java","lineNumber":2511,"sourceCode":"      public PTransform<PBegin, PCollection<Row>> buildExternal(\n          Read.External.Configuration config) {\n        Read.Builder<K, V> readBuilder = new AutoValue_KafkaIO_Read.Builder<>();\n        Read.Builder.setupExternalBuilder(readBuilder, config);\n\n        Class<Deserializer<K>> keyDeserializer =\n            (Class<Deserializer<K>>) resolveClass(config.keyDeserializer);\n        Coder<K> keyCoder = Read.Builder.resolveCoder(keyDeserializer);\n        if (!(keyCoder instanceof NullableCoder\n            && keyCoder.getCoderArguments().get(0) instanceof ByteArrayCoder)) {\n          throw new RuntimeException(\n              \"ExternalWithMetadata transform only supports keys of type nullable(byte[])\");\n        }\n        Class<Deserializer<V>> valueDeserializer =\n            (Class<Deserializer<V>>) resolveClass(config.valueDeserializer);\n        Coder<V> valueCoder = Read.Builder.resolveCoder(valueDeserializer);\n        if (!(valueCoder instanceof NullableCoder\n            && valueCoder.getCoderArguments().get(0) instanceof ByteArrayCoder)) {\n          throw new RuntimeException(\n              \"ExternalWithMetadata transform only supports values of type nullable(byte[])\");\n        }\n\n        return readBuilder.build().externalWithMetadata();\n      }\n    }\n\n    public static <K, V> ByteArrayKafkaRecord toExternalKafkaRecord(KafkaRecord<K, V> kafkaRecord) {\n      List<KafkaHeader> headers =\n          (kafkaRecord.getHeaders() == null)\n              ? null\n              : Arrays.stream(kafkaRecord.getHeaders().toArray())\n                  .map(h -> new KafkaHeader(h.key(), h.value()))\n                  .collect(Collectors.toList());\n      ByteArrayKafkaRecord byteArrayKafkaRecord =\n          new ByteArrayKafkaRecord(\n              kafkaRecord.getTopic(),\n              kafkaRecord.getPartition(),","sourceCodeStart":2493,"sourceCodeEnd":2529,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java#L2493-L2529","documentation":"Same restriction as for keys, applied to the value deserializer: the external-with-metadata KafkaIO transform only supports value coders of nullable(byte[]), since records are shipped as raw bytes plus metadata across the language boundary.","triggerScenarios":"Expanding KafkaIO externalWithMetadata with config.valueDeserializer that resolves to anything other than NullableCoder(ByteArrayCoder) — e.g. KafkaAvroDeserializer or StringDeserializer for values.","commonSituations":"Using Avro/Confluent deserializer in the cross-language read (should be handled after read via schema registry in the consuming transform), or StringDeserializer for values.","solutions":["Set valueDeserializer to org.apache.kafka.common.serialization.ByteArrayDeserializer.","Decode Avro/JSON payloads downstream using ParseAvro/ParseJson or a SchemaRegistry transform after the read.","Use plain Java KafkaIO.read() with explicit deserializers/coders when off the cross-language path."],"exampleFix":"// before\nconfig.valueDeserializer = KafkaAvroDeserializer.class.getName()\n// after\nconfig.valueDeserializer = \"org.apache.kafka.common.serialization.ByteArrayDeserializer\"","handlingStrategy":"validation","validationCode":"Coder<?> valueCoder = KafkaIO.Read.Builder.resolveCoder(valueDeserializerClass);\nif (!(valueCoder instanceof NullableCoder && valueCoder.getCoderArguments().get(0) instanceof ByteArrayCoder)) {\n  throw new IllegalArgumentException(\"External with-metadata read requires nullable(byte[]) values; parse Avro/JSON downstream\");\n}","typeGuard":"boolean isNullableBytes(Coder<?> c) {\n  return c instanceof NullableCoder\n      && !c.getCoderArguments().isEmpty()\n      && c.getCoderArguments().get(0) instanceof ByteArrayCoder;\n}","tryCatchPattern":"try {\n  pipeline.apply(KafkaIO.readAllExternalWithMetadata(config));\n} catch (RuntimeException e) {\n  if (e.getMessage() != null && e.getMessage().contains(\"only supports values\")) {\n    config.setValueDeserializer(\"org.apache.kafka.common.serialization.ByteArrayDeserializer\");\n  }\n}","preventionTips":["Never set Avro/Confluent deserializers in external-with-metadata configs","Do payload decoding (ParseAvro/SchemaRegistry transforms) after the read","Add a config sanity test asserting ByteArrayDeserializer for both key and value"],"tags":["java","kafka","cross-language","coder"],"backgroundTag":"type-mismatch","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"}