Sub-issue of #601.
Summary
Add a pluggable PayloadDeserializer interface and optional JSON Schema validation per topic, so that payloads are parsed, validated, and optionally rejected before reaching the SSE pipeline or MongoDB writer.
Motivation
The v1 module (#601) treats all payloads as opaque byte[]. In production, brokers emit payloads in various formats (JSON, plain text, binary Protobuf, Avro). Without a deserializer layer, every consumer must parse the raw bytes independently and handle malformed payloads individually.
Proposed design
PayloadDeserializer interface
public interface PayloadDeserializer<T> {
T deserialize(byte[] payload, String topic) throws PayloadDeserializationException;
}
Built-in implementations:
Utf8StringDeserializer (default — current behaviour)
JsonNodeDeserializer — parse as Jackson JsonNode; invalid JSON routes to dead-letter
ProtobufDeserializer — configurable .proto descriptor path
AvroDeserializer — configurable Avro schema file or Schema Registry URL
Custom deserializers registered as @RegisterPlugin-annotated classes.
JSON Schema validation
plugins-args:
mqtt-client:
subscriptions:
- topic: "sensors/#"
qos: 1
schema: "classpath:schemas/sensor-reading.json" # JSON Schema draft-07
on-validation-failure: "dead-letter" # "dead-letter" | "drop" | "pass"
Invalid payloads are routed to the dead-letter log (see #607) with the validation error attached.
Scope
PayloadDeserializer interface added to commons or as an internal mqtt SPI.
- Deserializer resolved per topic from config; falls back to
Utf8StringDeserializer.
- JSON Schema validation using
com.networknt:json-schema-validator.
- Unit tests for each built-in deserializer; integration test for validation + dead-letter routing.
Dependencies
Sub-issue of #601.
Summary
Add a pluggable
PayloadDeserializerinterface and optional JSON Schema validation per topic, so that payloads are parsed, validated, and optionally rejected before reaching the SSE pipeline or MongoDB writer.Motivation
The v1 module (#601) treats all payloads as opaque
byte[]. In production, brokers emit payloads in various formats (JSON, plain text, binary Protobuf, Avro). Without a deserializer layer, every consumer must parse the raw bytes independently and handle malformed payloads individually.Proposed design
PayloadDeserializer interface
Built-in implementations:
Utf8StringDeserializer(default — current behaviour)JsonNodeDeserializer— parse as JacksonJsonNode; invalid JSON routes to dead-letterProtobufDeserializer— configurable.protodescriptor pathAvroDeserializer— configurable Avro schema file or Schema Registry URLCustom deserializers registered as
@RegisterPlugin-annotated classes.JSON Schema validation
Invalid payloads are routed to the dead-letter log (see #607) with the validation error attached.
Scope
PayloadDeserializerinterface added tocommonsor as an internalmqttSPI.Utf8StringDeserializer.com.networknt:json-schema-validator.Dependencies
on-validation-failure: dead-letteris used.