This is an automated email from the ASF dual-hosted git repository.
riemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 50718cb4e9 fix: Support multiple sensors of the same type in OI4
adapter (#3026)
50718cb4e9 is described below
commit 50718cb4e9e961801c4349c448574744f8671985
Author: Dominik Riemer <[email protected]>
AuthorDate: Mon Jul 15 08:00:36 2024 +0200
fix: Support multiple sensors of the same type in OI4 adapter (#3026)
---
.../streampipes/connect/iiot/adapters/oi4/Oi4Adapter.java | 15 ++++++++++-----
1 file changed, 10 insertions(+), 5 deletions(-)
diff --git
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/oi4/Oi4Adapter.java
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/oi4/Oi4Adapter.java
index 06bd974c6c..aa24c3ec5d 100644
---
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/oi4/Oi4Adapter.java
+++
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/adapters/oi4/Oi4Adapter.java
@@ -149,7 +149,7 @@ public class Oi4Adapter implements StreamPipesAdapter {
InputStream in = convertByte(mqttEvent);
var networkMessage = mapper.readValue(in, NetworkMessage.class);
var payload = extractPayload(networkMessage);
- collector.collect(payload);
+ payload.forEach(collector::collect);
} catch (ParseException e) {
LOG.debug("Message parsing failed - this might be caused by
messages from a different sensor type");
} catch (IOException e) {
@@ -253,7 +253,7 @@ public class Oi4Adapter implements StreamPipesAdapter {
var networkMessage = mapper.readValue(sampleMessage,
NetworkMessage.class);
var payload = extractPayload(networkMessage);
- String plainPayload = mapper.writeValueAsString(payload);
+ String plainPayload = mapper.writeValueAsString(payload.get(0));
return new JsonParsers(new JsonObjectParser())
.getGuessSchema(convertByte(plainPayload.getBytes(StandardCharsets.UTF_8)));
} catch (IOException e) {
@@ -321,9 +321,10 @@ public class Oi4Adapter implements StreamPipesAdapter {
return new MqttConsumer(this.mqttConfig, eventProcessor);
}
- private Map<String, Object> extractPayload(NetworkMessage message) throws
ParseException {
+ private List<Map<String, Object>> extractPayload(NetworkMessage message)
throws ParseException {
var dataMessages = findProcessDataInputMessage(message);
+ var result = new ArrayList<Map<String, Object>>();
if (!dataMessages.isEmpty()) {
@@ -338,12 +339,16 @@ public class Oi4Adapter implements StreamPipesAdapter {
// an empty list of selected sensors means that we want to collect
data from all sensors available
if (selectedSensors.isEmpty() || selectedSensors.contains(sensorId))
{
- return extractAndEnrichMessagePayload(dataMessage, sensorId);
+ result.add(extractAndEnrichMessagePayload(dataMessage, sensorId));
}
}
}
}
- throw new ParseException(String.format("No sensor of type %s found in
message", givenSensorType));
+ if (!result.isEmpty()) {
+ return result;
+ } else {
+ throw new ParseException(String.format("No sensor of type %s found in
message", givenSensorType));
+ }
}
private List<DataSetMessage> findProcessDataInputMessage(NetworkMessage
message) {