Topic wildcards
Last updated
sensors/+/temperaturehumiditymqttConsumer:
broker: tcp://127.0.0.1:1883
topic: sensors/groundfloor/+/temperaturehumidity
clientId: tempHumidyManagementProcessor
qos: 1
deserializer:
parser: com.fractalworks.mqtt.example.TemperatureHumiditySensorParser
compressed: false
batch: falsepublic class TemperatureHumiditySensorParser implements StreamEventParser<List<StreamEvent>> {
private String ROOM = "room";
private ObjectMapper objectMapper;
public TemperatureHumiditySensorParser() {
// REQUIRED
}
@Override
public void initialise() throws StreamsException {
objectMapper = new ObjectMapper();
objectMapper.configure(SerializationFeature.INDENT_OUTPUT, false);
objectMapper.configure(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS, true);
objectMapper.setSerializationInclusion(JsonInclude.Include.NON_EMPTY);
}
@Override
public List<StreamEvent> translate(Object payload) throws TranslationException {
TemperatureHumidityEvent event = parseSingleEvent(payload);
if(event == null){
throw new TranslationException("Expected TemperatureHumidityEvent type to parse.");
}
return Collections.singletonList(event.toStreamEvent());
}
@Override
public List<StreamEvent> translate(Object payload, String source) throws TranslationException {
// First parse the event to the standard data type
List<StreamEvent> events = translate(payload);
// Now decorate the parsed events with room identifier
String[] tokens = source.split("/");
if(tokens.length != 4){
// Now decorate the event with the source information
events.forEach(event -> {
event.addValue(ROOM, tokens[3]);
});
}
return events;
}
/**
* Parse object to a single TemperatureHumidityEvent
*
* @param obj
* @return
*/
private TemperatureHumidityEvent parseSingleEvent(final Object obj) {
TemperatureHumidityEvent event = null;
try {
if (obj instanceof String str) {
event = objectMapper.readValue(str, TemperatureHumidityEvent.class);
} else if (obj instanceof byte[] byteArray) {
event = objectMapper.readValue(byteArray, TemperatureHumidityEvent.class);
}
} catch (IOException e) {
// Consume exception
}
return event;
}
/**
* Parse object to multiple stream events
*
* @param obj
* @return
*/
private List<TemperatureHumidityEvent> parseMultipleEvents(final Object obj) {
List<TemperatureHumidityEvent> events = null;
try {
if (obj instanceof String str) {
events = Arrays.asList(objectMapper.readValue(str, TemperatureHumidityEvent[].class));
} else if (obj instanceof byte[] byteArray) {
events = Arrays.asList(objectMapper.readValue(byteArray, TemperatureHumidityEvent[].class));
}
} catch (IOException e) {
// Consume exception
}
return events;
}
}// Some code
mqttConsumer:
broker: tcp://127.0.0.1:1883
topic: sensors/#
clientId: tempHumidyManagementProcessor
qos: 1
deserializer:
parser: com.fractalworks.mqtt.example.TemperatureHumiditySensorParser
compressed: false
batch: false