{"record":{"id":"915741b307f0dbef","repo":"apache/druid","slug":"kafka-deserializers-must-return-a-byte-array-byte","errorCode":null,"errorMessage":"Kafka deserializers must return a byte array (byte[]), %s returns %s","messagePattern":"Kafka deserializers must return a byte array \\(byte\\[\\]\\), (.+?) returns (.+?)","errorType":"exception","errorClass":"IllegalArgumentException","httpStatus":null,"severity":"error","filePath":"extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaRecordSupplier.java","lineNumber":357,"sourceCode":"    }\n  }\n\n  private static Deserializer getKafkaDeserializer(Properties properties, String kafkaConfigKey, boolean isKey)\n  {\n    Deserializer deserializerObject;\n    try {\n      Class deserializerClass = Class.forName(properties.getProperty(\n          kafkaConfigKey,\n          ByteArrayDeserializer.class.getTypeName()\n      ));\n      Method deserializerMethod = deserializerClass.getMethod(\"deserialize\", String.class, byte[].class);\n\n      Type deserializerReturnType = deserializerMethod.getGenericReturnType();\n\n      if (deserializerReturnType == byte[].class) {\n        deserializerObject = (Deserializer) deserializerClass.getConstructor().newInstance();\n      } else {\n        throw new IllegalArgumentException(\"Kafka deserializers must return a byte array (byte[]), \" +\n                                           deserializerClass.getName() + \" returns \" +\n                                           deserializerReturnType.getTypeName());\n      }\n    }\n    catch (ClassNotFoundException | NoSuchMethodException | InstantiationException | IllegalAccessException | InvocationTargetException e) {\n      throw new StreamException(e);\n    }\n\n    Map<String, Object> configs = new HashMap<>();\n    for (String key : properties.stringPropertyNames()) {\n      configs.put(key, properties.getProperty(key));\n    }\n\n    deserializerObject.configure(configs, isKey);\n    return deserializerObject;\n  }\n\n  public static KafkaConsumer<byte[], byte[]> getKafkaConsumer(","sourceCodeStart":339,"sourceCodeEnd":375,"githubUrl":"https://github.com/apache/druid/blob/9b90983fd291f26935af934383ce360473179e4d/extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaRecordSupplier.java#L339-L375","documentation":"KafkaRecordSupplier.getKafkaDeserializer() reflects on the configured deserializer class and requires its deserialize() method to return byte[]; Druid's kafka ingestion pipeline depends on raw bytes. A deserializer with a different return type throws IllegalArgumentException with the class name and actual return type.","triggerScenarios":"Setting kafka.consumer property key.deserializer or value.deserializer to a class whose deserialize method returns a non-byte[] type (e.g. String, custom object) instead of using the default ByteArrayDeserializer or returning bytes from a custom Deserializer.","commonSituations":"Pointing deserializer at a Kafka-provided typed deserializer like StringDeserializer or JsonSerializer inside Druid's consumer config.","solutions":["Use org.apache.kafka.common.serialization.ByteArrayDeserializer for key and value","If using a custom deserializer, make its deserialize() return byte[]","Do the decoding (String/JSON) inside Druid's input format (parser) instead of the Kafka deserializer"],"exampleFix":"// before\n\"key.deserializer\": \"org.apache.kafka.common.serialization.StringDeserializer\"\n// after\n\"key.deserializer\": \"org.apache.kafka.common.serialization.ByteArrayDeserializer\"","handlingStrategy":"validation","validationCode":"Class<?> c = Class.forName(deserializerClassName);\nMethod m = c.getMethod(\"deserialize\", String.class, byte[].class);\nif (!m.getReturnType().equals(byte[].class)) {\n  throw new IllegalArgumentException(deserializerClassName + \" must return byte[]\");\n}","typeGuard":null,"tryCatchPattern":"try {\n  supplier.start();\n} catch (IllegalArgumentException e) {\n  if (e.getMessage().startsWith(\"Kafka deserializers must return a byte array\")) {\n    // replace key/value.deserializer with ByteArrayDeserializer\n  } else { throw e; }\n}","preventionTips":["Always use ByteArrayDeserializer for Druid kafka ingestion","Do decoding in the input format, not the kafka deserializer","Validate consumer properties before task submission"],"tags":["kafka","deserializer","configuration"],"backgroundTag":"invalid-argument-value","analyzedSha":"9b90983fd291f26935af934383ce360473179e4d","analyzedAt":"2026-09-07T13:32:30.957Z","contentChangedAt":"2026-09-07T13:32:30.957Z","schemaVersion":2},"datasetVersion":"2026-09-17T15:17:12.973Z"}