apache/druid · error · IllegalArgumentException
Kafka deserializers must return a byte array (byte[])
Error message
Kafka deserializers must return a byte array (byte[]), %s returns %s
What it means
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.
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
Example fix
// before "key.deserializer": "org.apache.kafka.common.serialization.StringDeserializer" // after "key.deserializer": "org.apache.kafka.common.serialization.ByteArrayDeserializer"
Defensive patterns
Strategy: validation
Validate before calling
Class<?> c = Class.forName(deserializerClassName);
Method m = c.getMethod("deserialize", String.class, byte[].class);
if (!m.getReturnType().equals(byte[].class)) {
throw new IllegalArgumentException(deserializerClassName + " must return byte[]");
} Try / catch
try {
supplier.start();
} catch (IllegalArgumentException e) {
if (e.getMessage().startsWith("Kafka deserializers must return a byte array")) {
// replace key/value.deserializer with ByteArrayDeserializer
} else { throw e; }
} Prevention
- Always use ByteArrayDeserializer for Druid kafka ingestion
- Do decoding in the input format, not the kafka deserializer
- Validate consumer properties before task submission
When it happens
Trigger: 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.
Common situations: Pointing deserializer at a Kafka-provided typed deserializer like StringDeserializer or JsonSerializer inside Druid's consumer config.
Understand the failure class
Background: "Must be a positive integer", "Invalid value", "Unsupported": the invalid-argument-value error family, when a library rejects the value you pass — this error's family across 35 libraries.
Related errors
- No valid task counts after applying constraints for…
- Reset with skipped offsets is not supported when…
- Unable to create RecordSupplier
- A valid tlsPort needs to specified when druid.enableTlsPort…
- Already shut down, not starting again
AI-assisted analysis of apache/druid@9b90983fd2 (2026-09-07).
Data as JSON: /api/errors/915741b307f0dbef.
Report an issue: GitHub.
Appendix: source
Thrown at extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/KafkaRecordSupplier.java:357
}
}
private static Deserializer getKafkaDeserializer(Properties properties, String kafkaConfigKey, boolean isKey)
{
Deserializer deserializerObject;
try {
Class deserializerClass = Class.forName(properties.getProperty(
kafkaConfigKey,
ByteArrayDeserializer.class.getTypeName()
));
Method deserializerMethod = deserializerClass.getMethod("deserialize", String.class, byte[].class);
Type deserializerReturnType = deserializerMethod.getGenericReturnType();
if (deserializerReturnType == byte[].class) {
deserializerObject = (Deserializer) deserializerClass.getConstructor().newInstance();
} else {
throw new IllegalArgumentException("Kafka deserializers must return a byte array (byte[]), " +
deserializerClass.getName() + " returns " +
deserializerReturnType.getTypeName());
}
}
catch (ClassNotFoundException | NoSuchMethodException | InstantiationException | IllegalAccessException | InvocationTargetException e) {
throw new StreamException(e);
}
Map<String, Object> configs = new HashMap<>();
for (String key : properties.stringPropertyNames()) {
configs.put(key, properties.getProperty(key));
}
deserializerObject.configure(configs, isKey);
return deserializerObject;
}
public static KafkaConsumer<byte[], byte[]> getKafkaConsumer(View on GitHub (pinned to 9b90983fd2)