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) {

Reply via email to