This is an automated email from the ASF dual-hosted git repository.

zehnder pushed a commit to branch 
2964-timestamp-conversion-is-broken-in-adapters
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to 
refs/heads/2964-timestamp-conversion-is-broken-in-adapters by this push:
     new 644a319132 refactor(#2964): Refactor method onAdapterStarted in 
FileReplayAdapter
644a319132 is described below

commit 644a31913259dc44b4c15370ae1275ff07489726
Author: Philipp Zehnder <[email protected]>
AuthorDate: Mon Jul 1 15:20:49 2024 +0200

    refactor(#2964): Refactor method onAdapterStarted in FileReplayAdapter
---
 .../iiot/protocol/stream/FileReplayAdapter.java    | 55 ++++++++++++++--------
 1 file changed, 35 insertions(+), 20 deletions(-)

diff --git 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
index 6a8b14650f..cdc3fb371d 100644
--- 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
+++ 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/FileReplayAdapter.java
@@ -135,8 +135,39 @@ public class FileReplayAdapter implements 
StreamPipesAdapter {
       IAdapterRuntimeContext adapterRuntimeContext
   ) throws AdapterException {
 
-    // extract user input
+    boolean replayOnce = 
extractUserInputsAndReturnValueOfReplayOnce(extractor);
+
+    determineTimestampRuntimeName(extractor);
+
+    determineSourceTimestampField(extractor);
+
+    startAdapterReplayThread(extractor, collector, adapterRuntimeContext, 
replayOnce);
+  }
+
+  private void startAdapterReplayThread(
+      IAdapterParameterExtractor extractor,
+      IEventCollector collector,
+      IAdapterRuntimeContext adapterRuntimeContext,
+      boolean replayOnce
+  ) {
     executor = Executors.newScheduledThreadPool(1);
+    if (replayOnce) {
+      executor.schedule(
+          () -> getFileFromEndpointAndParseFile(extractor, collector, 
adapterRuntimeContext),
+          0,
+          TimeUnit.SECONDS
+      );
+    } else {
+      executor.scheduleAtFixedRate(
+          () -> getFileFromEndpointAndParseFile(extractor, collector, 
adapterRuntimeContext),
+          0,
+          1,
+          TimeUnit.SECONDS
+      );
+    }
+  }
+
+  private boolean 
extractUserInputsAndReturnValueOfReplayOnce(IAdapterParameterExtractor 
extractor) {
     boolean replayOnce = extractor
         .getStaticPropertyExtractor()
         .selectedSingleValue(REPLAY_ONCE, String.class)
@@ -155,8 +186,10 @@ public class FileReplayAdapter implements 
StreamPipesAdapter {
           .singleValueParameter(SPEED_UP, Float.class);
       default -> 1.0f;
     };
+    return replayOnce;
+  }
 
-    // get timestamp field
+  private void determineTimestampRuntimeName(IAdapterParameterExtractor 
extractor) throws AdapterException {
     var timestampField = extractor
         .getAdapterDescription()
         .getEventSchema()
@@ -173,24 +206,6 @@ public class FileReplayAdapter implements 
StreamPipesAdapter {
       timestampRuntimeName = timestampField.get()
                                            .getRuntimeName();
     }
-
-    determineSourceTimestampField(extractor);
-
-    // start replay adapter
-    if (replayOnce) {
-      executor.schedule(
-          () -> getFileFromEndpointAndParseFile(extractor, collector, 
adapterRuntimeContext),
-          0,
-          TimeUnit.SECONDS
-      );
-    } else {
-      executor.scheduleAtFixedRate(
-          () -> getFileFromEndpointAndParseFile(extractor, collector, 
adapterRuntimeContext),
-          0,
-          1,
-          TimeUnit.SECONDS
-      );
-    }
   }
 
   /**

Reply via email to